---
title: "01-Elasticsearch 学习笔记"
created: 2025-12-15
aliases:
- Elasticsearch 学习笔记
tags:
- 项目
---
# Elasticsearch 学习笔记
> Elasticsearch 写入时先将每篇文章分词,再反向建立 “单个词汇→包含该词汇的所有文章” 的倒排索引,同时对词汇排序以支撑后续高效查询;搜索时则先借助内存中的 Term Index 前缀定位 + 二分查找快速找到目标词汇对应的文章 ID 列表,再根据 AND/OR 等搜索条件对多个词汇的文章 ID 列表取交集或并集,最后依据筛选后的文章 ID 提取完整内容并返回。
[B站视频 BV1yb421J7oX](https://player.bilibili.com/player.html?bvid=BV1yb421J7oX&high_quality=1&autoplay=0#video-iframe)
## **一、顶层设计:为什么需要 Elasticsearch?**
### **1.1 ES 的定位**
![[2-Learning/05-项目/08-企业级项目深读/02-damai_pro/05-业务深读/04-节目服务深度解析/assets/Elasticsearch_在系统中的定位-9b2df756.jpg]]
### **1.2 传统数据库面临的问题**
![[2-Learning/05-项目/08-企业级项目深读/02-damai_pro/05-业务深读/04-节目服务深度解析/assets/传统关系型数据库的瓶颈-2e271bb1.jpg]]
- **模糊查询性能瓶颈**:MySQL 等关系型数据库使用 `LIKE %keyword%` 进行左模糊或全模糊查询时,无法利用索引(Index),会导致**全表扫描**,在百万级数据量下性能急剧下降。
- **复杂维度筛选困难**:当面临海量数据的多维度筛选、统计、标签聚合时,SQL 语句变得极度复杂且执行效率低下。
- **文本相关性缺失**:传统数据库只能做“精准匹配”或“简单的包含匹配”,无法计算**相关性得分(Score)**,无法按“搜索结果匹配度”对结果进行排序。
- **分词能力弱**:无法处理同义词、纠错、中文分词(如将“由于”和“游泳”区分开)等复杂的自然语言处理需求。
### **1.3 Elasticsearch 解决的核心问题**
![[2-Learning/05-项目/08-企业级项目深读/02-damai_pro/05-业务深读/04-节目服务深度解析/assets/Elasticsearch_核心价值-0a0cdb5f.jpg]]
- **性能全文检索**:基于**倒排索引**,实现亿级数据毫秒级响应。
- **高可用与横向扩展**:天生的分布式架构,通过增加节点即可线性提升存储容量和计算能力。
- **复杂聚合分析**:提供强大的聚合(Aggregations)框架,能够替代部分 OLAP 场景,实时统计数据分布。
- **相关性排序**:基于 TF-IDF / BM25 算法,提供符合人类直觉的搜索结果排名。
## **二、核心原理:如何解决问题?**
### **2.1 倒排索引(Inverted Index)**
这是 ES 快如闪电的根本原因。
- **正排索引(Forward Index)**:文档 ID -> 文档内容(类似 MySQL 主键查询)。
- **倒排索引**:单词(Term)-> 包含该单词的文档 ID 列表。
- ES 在写入数据时,会通过**分词器(Analyzer)**将文本拆解为单词,建立索引。
- 查询时,直接根据单词找到文档 ID 列表,并通过位运算快速合并结果,无需扫描全表。
![[2-Learning/05-项目/08-企业级项目深读/02-damai_pro/05-业务深读/04-节目服务深度解析/assets/倒排索引原理详解-5d950c3c.jpg]]
### **2.2 Term Index(词项索引)**
![[2-Learning/05-项目/08-企业级项目深读/02-damai_pro/05-业务深读/04-节目服务深度解析/assets/Term_Index_加速原理-fe766372.jpg]]
- **面临的问题**:Term Dictionary(词项字典)数据量极大(千万/亿级),无法全部放入内存,必须存在磁盘。每次查询都去磁盘翻字典,I/O 消耗巨大。
- **解决方案**:**Term Index** 是一个基于词项前缀构建的**精简目录树**(类似 Trie 树或 FST - Finite State Transducers)。
- **内存驻留**:它体积非常小,可以完全加载到 **内存(RAM)** 中。
- **加速原理**:查询时,先在内存中查 Term Index,找到该 Term 在磁盘 Term Dictionary 中的大概位置(Offset 偏移量),然后再去磁盘读取具体内容。
- *类比*:Term Dictionary 是字典本身(在磁盘),Term Index 是字典的“拼音首字母索引页”(在内存)。
### **2.3 存储结构**
![[2-Learning/05-项目/08-企业级项目深读/02-damai_pro/05-业务深读/04-节目服务深度解析/assets/ES_存储结构详解-c3bd152d.jpg]]
当倒排索引帮我们找到文档 ID 后,我们还需要获取内容或进行排序。
- **Stored Fields(行式存储)**
- **用途**:用于存储**原始文档内容**(`_source` 字段)。
- **特点**:行式存储,适合展示完整信息,但不适合聚合分析(读取很多冗余数据)。
- **Doc Values(列式存储)**
- **用途**:专用于**排序(Sorting)和聚合(Aggregations)**。
- **原理**:空间换时间。将散落在不同文档中的同一个字段值,集中在一起列式存储。
- **优势**:排序或计算平均值时,CPU 可以连续读取内存地址,极大提升效率,且对操作系统文件缓存(OS Cache)非常友好。
### **2.4 Segment 与 Lucene**
![[2-Learning/05-项目/08-企业级项目深读/02-damai_pro/05-业务深读/04-节目服务深度解析/assets/Segment_与_Lucene_架构-6f82b8e4.jpg]]
- **Lucene**:ES 的底层核心库。一个 ES 的分片(Shard)本质上就是一个完整的 **Lucene 索引**。
- **Segment(段)**:
- **最小单元**:Lucene 内部由多个 Segment 组成。一个 Segment 包含了倒排索引、Term Index、Stored Fields 等所有结构,是具备完整搜索功能的最小单元。
- **不可变性(Immutable)**:Segment 一旦生成,**不可修改**。
- **写入与合并**:新增数据会生成新的 Segment;删除数据只是打上 `.del` 标记(逻辑删除)。ES 会在后台自动进行 **Segment Merging(段合并)**,将多个小 Segment 合并为大 Segment,同时物理剔除被标记删除的数据。
## **三、分布式架构设计**
![[2-Learning/05-项目/08-企业级项目深读/02-damai_pro/05-业务深读/04-节目服务深度解析/assets/ES_分布式架构设计-01752588.jpg]]
- **Cluster(集群)**:由多个节点组成。
- **Node(节点)**:单个服务实例。
- **Index(索引)**:逻辑上的数据集合(类比 Database)。
- **Shard(分片)**:数据被切分成多个分片存储在不同节点,实现并行计算和存储扩展。
- **Replica(副本)**:分片的备份,用于提高可用性(HA)和读取吞吐量。
### **3.1 高性能优化:分片(Shard)**
![[2-Learning/05-项目/08-企业级项目深读/02-damai_pro/05-业务深读/04-节目服务深度解析/assets/高性能优化:分片机制-6bcd256e.jpg]]
- **问题**:单个 Index 数据量过大(如 1TB),单机硬盘存不下,且搜索时单线程扫描太慢。
- **解决方案**:**分片机制**。
- 将一个 Index 逻辑拆分为多个 Shard(分片)。
- 每个 Shard 是一个独立的 Lucene 实例。
- **优势**:读写压力被分散到多个 Shard 上并行处理,极大提升吞吐量。
### **3.2 高扩展性优化:多节点(Node)**
![[2-Learning/05-项目/08-企业级项目深读/02-damai_pro/05-业务深读/04-节目服务深度解析/assets/高扩展性优化:多节点部署-fa229e0c.jpg]]
- **原理**:**横向扩展(Scale Out)**。
- **机制**:当数据量增长,分片增多,单机 CPU/内存吃紧时,可以向集群中加入新的机器(Node)。
- **自动平衡**:ES 会自动感知新节点,并将部分 Shard 迁移过去,实现负载均衡。
### **3.3 高可用优化:副本(Replica)**
![[2-Learning/05-项目/08-企业级项目深读/02-damai_pro/05-业务深读/04-节目服务深度解析/assets/高可用优化:副本机制-27535f46.jpg]]
- **角色区分**:
- **Primary Shard(主分片)**:负责处理写入请求。
- **Replica Shard(副本分片)**:主分片的完整备份。
- **机制**:
- **读写分离**:副本分片可以分担搜索(读)请求,提升查询并发量。
- **故障转移(Failover)**:如果持有主分片的 Node 挂掉,集群会迅速选举一个副本分片升级为新的主分片,保证服务不中断。
### **3.4 节点角色分化**
![[2-Learning/05-项目/08-企业级项目深读/02-damai_pro/05-业务深读/04-节目服务深度解析/assets/Node_角色分化-523fd620.jpg]]
在大型集群中,让每个节点“各司其职”效率更高。
- **Master Node(主节点)**:集群的大脑。负责索引创建/删除、维护集群状态(Cluster State)、管理节点加入/退出。
- **Data Node(数据节点)**:苦力。负责存储数据(Shard),执行耗费资源的 CRUD 和聚合操作。
- **Coordinate Node(协调节点)**:前台/路由。接收客户端请求,分发给 Data Node,并汇总最终结果(Scatter-Gather)。
### **3.5 去中心化协调机制**
![[2-Learning/05-项目/08-企业级项目深读/02-damai_pro/05-业务深读/04-节目服务深度解析/assets/去中心化协调机制-0c561c9f.jpg]]
- **Raft 算法衍生**:ES 内部实现了一套基于 Raft 改进的共识算法(Zen Discovery)。
- **作用**:
- 保证集群中所有节点对“集群状态”的认知是一致的。
- 实现 Master 节点的选举。
- **故障检测**:节点之间互相 Ping,感知是否有节点掉线。
## **四、核心流程详解**
### **4.1 写入流程**
ES 的写入流程设计权衡了**数据安全性**与**写入高吞吐**。
**1. 路由与转发**
- 客户端向任意节点发送写入请求(该节点暂时成为**协调节点**)。
- 协调节点使用**路由算法**确定数据所属的主分片(Primary Shard)位置:
- `shard = hash(routing) % number_of_primary_shards`
- *注:*`routing` *默认是文档* `_id`*。*
- 协调节点将请求转发给持有该主分片的 Data Node。
**2. 主分片写入 (Primary Operation)**
- **写入内存缓冲区 (Memory Buffer)**:数据先写入内存 buffer,此时数据**不可被搜索**。
- **写入 Translog (Transaction Log)**:同时追加写入 Translog 文件(顺序写磁盘),防止断电丢失数据。
- **Refresh (准实时关键步骤)**:默认每 1 秒,ES 将 buffer 中的数据生成一个新的 **Segment** 文件(此时建立倒排索引),并清空 buffer。**一旦生成 Segment,数据即可被搜索**。这就是 ES 被称为“准实时(Near Real-Time, NRT)”的原因。
**3. 同步副本 (Replication)**
- 主分片写入成功后,并行将请求发送给所有的 **副本分片 (Replica Shards)**。
- 副本分片执行相同的写入逻辑。
**4. 响应客户端**
- 当所有在 **ISR (In-Sync Replicas)** 列表中的副本都反馈写入成功后,主分片向协调节点报告成功。
- 协调节点向客户端返回“写入完成”。
> **技术深挖:Flush 操作** Translog 不会无限增长。当 Translog 达到阈值或每隔 30 分钟,ES 会触发 **Flush** 操作:
>
> 1. 强制执行 Refresh。
> 2. 将所有内存中的 Segment 强制 `fsync` 刷入物理磁盘。
> 3. 清空 Translog。 *这保证了数据的持久化存储。*
![[2-Learning/05-项目/08-企业级项目深读/02-damai_pro/05-业务深读/04-节目服务深度解析/assets/ES_写入流程详解-4fcf1ad9.jpg]]
### **4.2 搜索流程(Query Then Fetch)**
![[2-Learning/05-项目/08-企业级项目深读/02-damai_pro/05-业务深读/04-节目服务深度解析/assets/ES_搜索流程详解-f639461c.jpg]]
搜索比写入复杂,因为数据分散在多个分片上,必须通过“两阶段”策略来整合结果,以避免网络带宽的巨大浪费。
**阶段一:查询阶段 (Query Phase)**
- **请求分发**:客户端向**协调节点**发送搜索请求。协调节点根据请求(是否有 routing 参数)将请求广播到所有相关分片(主分片或副本分片均可,负载均衡)。
- **本地检索**:每个分片在本地 Lucene 中执行搜索:
1. 利用倒排索引筛选匹配文档。
2. 利用 Doc Values 进行排序和打分。
3. **关键点**:分片**仅返回**文档 ID、相关性算分 (\_score) 和排序值给协调节点,**不返回**文档的完整内容 (`_source`)。
- **全局排序**:协调节点收到所有分片返回的轻量级列表(例如每分片前 10 条),在内存中进行**全局归并排序**,选出最终的 Top N 文档 ID。
**阶段二:获取阶段 (Fetch Phase)**
- **精确定位**:协调节点知道了最终需要哪几个文档,以及它们位于哪个分片。
- **抓取内容**:协调节点向相关分片发送 `Multi-Get` 请求,只索取这 Top N 文档的完整内容 (`_source` / `Stored Fields`)。
- **返回结果**:分片返回文档详情,协调节点拼装最终 JSON,响应给客户端。
> **性能隐患:深度分页 (Deep Pagination)** 如果查询 `from=10000, size=10`:
>
> - 每个分片都必须查询出前 10010 条记录。
> - 假设有 5 个分片,协调节点需要接收 `5 * 10010 = 50050` 条记录的 ID,并在内存中排序,最后只取 10 条。
> - **后果**:内存爆炸,CPU 飙升。
> - **对策**:避免深分页,使用 `Search After` 或 `Scroll` API。
### **4.3 搜索流程总结图**
![[2-Learning/05-项目/08-企业级项目深读/02-damai_pro/05-业务深读/04-节目服务深度解析/assets/搜索流程数据结构使用-f707625b.jpg]]
## **五、核心概念对比与 Type 演变**
### **5.1 核心概念对比**
```text
┌─────────────────────────────────────────────────────────────────────────┐
│ ES 与关系型数据库概念对比 │
├───────────────────┬─────────────────────┬───────────────────────────────┤
│ 关系型数据库 │ Elasticsearch │ 说明 │
├───────────────────┼─────────────────────┼───────────────────────────────┤
│ Database │ Cluster │ 数据库/集群 │
│ Table │ Index │ 表/索引 │
│ Row │ Document │ 行/文档 │
│ Column │ Field │ 列/字段 │
│ Schema │ Mapping │ 表结构/映射 │
│ Index │ Inverted Index │ 索引/倒排索引 │
│ SQL │ Query DSL │ 查询语言 │
└───────────────────┴─────────────────────┴───────────────────────────────┘
```
### **5.2 Type 演变历史**
![[2-Learning/05-项目/08-企业级项目深读/02-damai_pro/05-业务深读/04-节目服务深度解析/assets/Type_概念的演变-da19af3e.jpg]]
- **5.x 及以前**:允许一个 Index 下存在多个 Type(类比 Table),但本质上底层字段是扁平化混在一起的,导致数据稀疏(Sparse)问题,影响压缩效率和性能。
- **6.x**:强制规定一个 Index 只能有一个 Type,通常默认名为 `doc`。
- **7.x**:Type 概念被彻底废弃(默认为 `_doc`),API 中 URL 的 type 参数变为可选。
- **8.x**:彻底移除 Type 概念。
- **结论**:现在设计索引时,严格遵循 **“一个 Index 对应一类业务数据”** 的原则。
**为什么移除?**
- 映射爆炸(Mapping Explosion)
- 在使用多类型时,如果不同类型之间有大量不同的字段,这会导致映射的数量急剧增加,进而引发映射爆炸问题。映射爆炸不仅会消耗大量的内存资源,还会降低 Elasticsearch 的性能,尤其是在处理大量数据时
- 字段名冲突
- 在同一个index的不同type中,如果有相同名称但映射类型不同的字段,会造成字段名冲突。这是因为 Elasticsearch 在内部是将这些字段扁平化处理的,而不同类型的相同名称字段可能会导致数据解析和查询时的混乱
- 我们可以和关系型数据库来对比,在同一个数据库中,这些不同的表,可以有名称相同但类型不同的字段。而在 Elasticsearch 同一个index的不同type中,如果有不同document的字段名相同,但是类型不同,就会报错
- 综上所述,Elasticsearch 从 7.x 版本开始废弃类型的主要目的是为了提升系统的性能、避免映射爆炸和字段冲突的问题,以及简化数据模型的设计和管理。这一改变反映了 Elasticsearch 对于提高性能、可维护性和用户体验的持续追求
## **六、核心功能**
### **6.1 搜索能力矩阵**
![[2-Learning/05-项目/08-企业级项目深读/02-damai_pro/05-业务深读/04-节目服务深度解析/assets/ES_搜索能力矩阵-cd6c4c21.jpg]]
### **6.2 聚合分析**
![[2-Learning/05-项目/08-企业级项目深读/02-damai_pro/05-业务深读/04-节目服务深度解析/assets/聚合分析类型2-1609a946.jpg]]
## **七、应用场景**
![[2-Learning/05-项目/08-企业级项目深读/02-damai_pro/05-业务深读/04-节目服务深度解析/assets/Elasticsearch_应用场景-39904ef5.jpg]]
1. **日志与监控(ELK Stack)**:收集服务器日志、应用 Error 日志,快速定位故障(Logstash/Beats + ES + Kibana)。
2. **站内搜索**:电商商品搜索、论坛帖子搜索、企业知识库检索。
3. **大屏可视化/BI**:实时统计大盘数据(如双11大屏),利用聚合功能快速出报表。
4. **地理位置服务(LBS)**:查询“附近的酒店”、“方圆5公里内的订单”(Geo-point/Geo-shape)。
## **八、如何使用 Elasticsearch**
### **8.1 Spring Boot 集成**
通常有两种主流方式:
1. **Spring Data Elasticsearch**:封装程度极高,类似 JPA/MyBatis-Plus,通过 Repository 接口操作。
- *优点*:开发极快,代码简洁。
- *缺点*:灵活性稍差,对复杂 DSL 和版本兼容性控制不如原生客户端细致。
2. **RestHighLevelClient (官方推荐/传统)**:基于 HTTP 的原生客户端封装。
- *优点*:完全覆盖官方 API,灵活,可控性强。
- *注意*:ES 7.15+ 后官方推出了新的 `Elasticsearch Java API Client`,但 `RestHighLevelClient` 依然在存量系统中广泛使用。**本文基于此方案进行封装。**
```xml
org.springframework.boot
spring-boot-starter-data-elasticsearch
```
```yaml
# application.yml 配置
spring:
elasticsearch:
uris: http://localhost:9200
username: elastic
password: password
```
此type的作用就是为了兼容6.x、7.x中的type概念,默认是关闭
### **8.2 原生 API 操作的复杂性**
直接使用 `RestHighLevelClient` 会面临大量样板代码:
- 构建 `SearchSourceBuilder`、`BoolQueryBuilder` 极其繁琐。
- 需要手动处理 `IOException`。
- 响应结果解析(Parse)需要从 JSON 层层剥离,非常痛苦。
- 连接管理和配置分散。
```java
// 原生 ES API 操作示例 - 复杂且繁琐
public SearchResponse searchPrograms(String keyword, Integer categoryId) {
// 构建查询条件
BoolQueryBuilder boolQuery = QueryBuilders.boolQuery();
if (StringUtils.isNotBlank(keyword)) {
boolQuery.must(QueryBuilders.matchQuery("title", keyword));
}
if (categoryId != null) {
boolQuery.filter(QueryBuilders.termQuery("categoryId", categoryId));
}
// 构建排序
FieldSortBuilder sortBuilder = SortBuilders.fieldSort("showTime")
.order(SortOrder.ASC);
// 构建搜索源
SearchSourceBuilder sourceBuilder = new SearchSourceBuilder();
sourceBuilder.query(boolQuery);
sourceBuilder.sort(sortBuilder);
sourceBuilder.from(0);
sourceBuilder.size(10);
// 构建搜索请求
SearchRequest searchRequest = new SearchRequest("program-index");
searchRequest.source(sourceBuilder);
// 执行搜索
return restHighLevelClient.search(searchRequest, RequestOptions.DEFAULT);
}
```
在平时开发中还是使用关系型数据库更加普遍,对于数据库、表、字段的概念更为熟悉,也更加习惯对表概念的操作。
而在操作Elasticsearch时,提供的api其实是很复杂的,各种操作的对象,如`SearchSourceBuilder FieldSortBuilder` `BoolQueryBuilder` 等等,操作上其实算不上简单,为了解决这个问题,在springboot操作Elasticsearch的基础上,进一步的封装,使用起来贴近于关系型数据库的方式,操作起来更加的容易上手
## **九、封装设计**
### **9.1 封装目标**
- **统一配置**:简化连接参数管理。
- **屏蔽细节**:隐藏繁琐的 Builder 构建过程。
- **简化查询**:通过 Map 或对象传递参数,自动构建 DSL。
- **结果转换**:自动将 ES 的 JSON 结果转为 Java Bean。
- **健壮性**:统一异常处理,防止 ES 波动导致服务崩溃。
![[2-Learning/05-项目/08-企业级项目深读/02-damai_pro/05-业务深读/04-节目服务深度解析/assets/封装设计目标-c8a9c448.jpg]]
### **9.2 封装架构设计**
![[2-Learning/05-项目/08-企业级项目深读/02-damai_pro/05-业务深读/04-节目服务深度解析/assets/ES_封装框架架构-9d1b66fd.jpg]]
```text
elasticsearch-spring-boot-starter/
├── src/main/java/com/example/elasticsearch/
│ ├── config/
│ │ ├── ElasticsearchProperties.java # 配置属性
│ │ └── ElasticsearchAutoConfiguration.java # 自动配置
│ ├── core/
│ │ ├── ElasticsearchService.java # 核心服务类
│ │ ├── ElasticsearchIndexService.java # 索引操作服务
│ │ └── ElasticsearchDocumentService.java # 文档操作服务
│ ├── query/
│ │ ├── EsQueryBuilder.java # 查询构建器
│ │ ├── EsSearchRequest.java # 搜索请求封装
│ │ └── EsHighlightConfig.java # 高亮配置
│ ├── result/
│ │ ├── PageResult.java # 分页结果
│ │ ├── SearchResult.java # 搜索结果
│ │ └── AggregationResult.java # 聚合结果
│ ├── exception/
│ │ └── ElasticsearchException.java # 自定义异常
│ └── annotation/
│ ├── EsDocument.java # 文档注解
│ └── EsField.java # 字段注解
└── src/main/resources/
└── META-INF/spring.factories
```
### **9.3 配置类设计**
```java
package com.example.elasticsearch.config;
import lombok.Data;
import org.springframework.boot.context.properties.ConfigurationProperties;
import org.springframework.validation.annotation.Validated;
import javax.validation.constraints.Min;
import javax.validation.constraints.NotEmpty;
import java.util.List;
/**
* Elasticsearch 配置属性类
* 支持集群配置、连接池、认证、重试等
*/
@Data
@Validated
@ConfigurationProperties(prefix = "elasticsearch")
public class ElasticsearchProperties {
/**
* 是否启用 ES
*/
private Boolean enabled = true;
/**
* ES 节点地址列表(支持集群)
*/
@NotEmpty(message = "ES 节点地址不能为空")
private List nodes = List.of("localhost:9200");
/**
* 用户名(可选)
*/
private String username;
/**
* 密码(可选)
*/
private String password;
/**
* 协议:http 或 https
*/
private String scheme = "http";
/**
* 是否启用 Type(兼容 ES 6.x 版本)
*/
private Boolean enableType = false;
/**
* 默认 Type 名称
*/
private String defaultType = "_doc";
/**
* 连接超时时间(毫秒)
*/
@Min(value = 1000, message = "连接超时时间不能小于1000ms")
private Integer connectTimeout = 5000;
/**
* Socket 超时时间(毫秒)
*/
@Min(value = 1000, message = "Socket超时时间不能小于1000ms")
private Integer socketTimeout = 30000;
/**
* 请求超时时间(毫秒)
*/
private Integer connectionRequestTimeout = 5000;
/**
* 最大连接数
*/
@Min(value = 1, message = "最大连接数不能小于1")
private Integer maxConnTotal = 100;
/**
* 每个路由的最大连接数
*/
@Min(value = 1, message = "每个路由最大连接数不能小于1")
private Integer maxConnPerRoute = 50;
/**
* 重试次数
*/
@Min(value = 0, message = "重试次数不能为负数")
private Integer retryTimes = 3;
/**
* 重试间隔(毫秒)
*/
private Long retryInterval = 1000L;
/**
* 是否开启嗅探器
*/
private Boolean enableSniffer = false;
/**
* 嗅探间隔时间(毫秒)
*/
private Long snifferInterval = 60000L;
/**
* 批量操作每批大小
*/
@Min(value = 100, message = "批量操作每批大小不能小于100")
private Integer bulkBatchSize = 1000;
/**
* 批量操作刷新策略:immediate, wait_for, none
*/
private String bulkRefreshPolicy = "none";
/**
* 是否打印 DSL 日志
*/
private Boolean printDsl = false;
/**
* 慢查询阈值(毫秒),超过此值记录警告日志
*/
private Long slowQueryThreshold = 3000L;
}
```
### **9.4 自动配置类**
```java
package com.example.elasticsearch.config;
import com.example.elasticsearch.core.ElasticsearchDocumentService;
import com.example.elasticsearch.core.ElasticsearchIndexService;
import com.example.elasticsearch.core.ElasticsearchService;
import lombok.extern.slf4j.Slf4j;
import org.apache.http.HttpHost;
import org.apache.http.auth.AuthScope;
import org.apache.http.auth.UsernamePasswordCredentials;
import org.apache.http.client.CredentialsProvider;
import org.apache.http.impl.client.BasicCredentialsProvider;
import org.elasticsearch.client.RestClient;
import org.elasticsearch.client.RestClientBuilder;
import org.elasticsearch.client.RestHighLevelClient;
import org.elasticsearch.client.sniff.Sniffer;
import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean;
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
import org.springframework.boot.context.properties.EnableConfigurationProperties;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.util.StringUtils;
import javax.annotation.PreDestroy;
import java.util.List;
import java.util.stream.Collectors;
/**
* Elasticsearch 自动配置类
*/
@Slf4j
@Configuration
@EnableConfigurationProperties(ElasticsearchProperties.class)
@ConditionalOnProperty(prefix = "elasticsearch", name = "enabled", havingValue = "true", matchIfMissing = true)
public class ElasticsearchAutoConfiguration {
private RestHighLevelClient restHighLevelClient;
private Sniffer sniffer;
@Bean
@ConditionalOnMissingBean
public RestHighLevelClient restHighLevelClient(ElasticsearchProperties properties) {
// 解析节点地址
List httpHosts = properties.getNodes().stream()
.map(node -> {
String[] parts = node.split(":");
String host = parts[0];
int port = parts.length > 1 ? Integer.parseInt(parts[1]) : 9200;
return new HttpHost(host, port, properties.getScheme());
})
.collect(Collectors.toList());
RestClientBuilder builder = RestClient.builder(
httpHosts.toArray(new HttpHost[0])
);
// 设置请求配置
builder.setRequestConfigCallback(requestConfigBuilder ->
requestConfigBuilder
.setConnectTimeout(properties.getConnectTimeout())
.setSocketTimeout(properties.getSocketTimeout())
.setConnectionRequestTimeout(properties.getConnectionRequestTimeout())
);
// 设置 HTTP 客户端配置
builder.setHttpClientConfigCallback(httpClientBuilder -> {
// 设置连接池
httpClientBuilder.setMaxConnTotal(properties.getMaxConnTotal());
httpClientBuilder.setMaxConnPerRoute(properties.getMaxConnPerRoute());
// 设置认证
if (StringUtils.hasText(properties.getUsername())
&& StringUtils.hasText(properties.getPassword())) {
CredentialsProvider credentialsProvider = new BasicCredentialsProvider();
credentialsProvider.setCredentials(
AuthScope.ANY,
new UsernamePasswordCredentials(
properties.getUsername(),
properties.getPassword()
)
);
httpClientBuilder.setDefaultCredentialsProvider(credentialsProvider);
}
return httpClientBuilder;
});
// 设置失败重试策略
builder.setFailureListener(new RestClient.FailureListener() {
@Override
public void onFailure(org.elasticsearch.client.Node node) {
log.warn("ES 节点 [{}] 连接失败", node.getHost());
}
});
restHighLevelClient = new RestHighLevelClient(builder);
// 启用嗅探器
if (properties.getEnableSniffer()) {
sniffer = Sniffer.builder(restHighLevelClient.getLowLevelClient())
.setSniffIntervalMillis(properties.getSnifferInterval().intValue())
.build();
log.info("ES 嗅探器已启用,间隔: {}ms", properties.getSnifferInterval());
}
log.info("ES 客户端初始化成功, 节点: {}", properties.getNodes());
return restHighLevelClient;
}
@Bean
@ConditionalOnMissingBean
public ElasticsearchService elasticsearchService(
RestHighLevelClient client,
ElasticsearchProperties properties) {
return new ElasticsearchService(client, properties);
}
@Bean
@ConditionalOnMissingBean
public ElasticsearchIndexService elasticsearchIndexService(
RestHighLevelClient client,
ElasticsearchProperties properties) {
return new ElasticsearchIndexService(client, properties);
}
@Bean
@ConditionalOnMissingBean
public ElasticsearchDocumentService elasticsearchDocumentService(
RestHighLevelClient client,
ElasticsearchProperties properties) {
return new ElasticsearchDocumentService(client, properties);
}
@PreDestroy
public void destroy() {
try {
if (sniffer != null) {
sniffer.close();
log.info("ES 嗅探器已关闭");
}
if (restHighLevelClient != null) {
restHighLevelClient.close();
log.info("ES 客户端已关闭");
}
} catch (Exception e) {
log.error("关闭 ES 客户端失败", e);
}
}
}
```
### **9.5 自定义异常类**
```java
package com.example.elasticsearch.exception;
import lombok.Getter;
/**
* Elasticsearch 自定义异常
*/
@Getter
public class ElasticsearchException extends RuntimeException {
private static final long serialVersionUID = 1L;
/**
* 索引名称
*/
private String index;
/**
* 操作类型
*/
private String operation;
/**
* 错误码
*/
private String errorCode;
/**
* 是否可重试
*/
private boolean retryable;
public ElasticsearchException(String message) {
super(message);
this.retryable = false;
}
public ElasticsearchException(String message, Throwable cause) {
super(message, cause);
this.retryable = isRetryableException(cause);
}
public ElasticsearchException(String operation, String index, String message) {
super(String.format("[%s] 索引 [%s] 操作失败: %s", operation, index, message));
this.operation = operation;
this.index = index;
this.errorCode = operation + "_ERROR";
}
public ElasticsearchException(String operation, String index, Throwable cause) {
super(String.format("[%s] 索引 [%s] 操作失败: %s",
operation, index, cause.getMessage()), cause);
this.operation = operation;
this.index = index;
this.errorCode = operation + "_ERROR";
this.retryable = isRetryableException(cause);
}
/**
* 判断异常是否可重试
*/
private boolean isRetryableException(Throwable cause) {
if (cause == null) {
return false;
}
String message = cause.getMessage();
if (message == null) {
return false;
}
// 可重试的异常类型
return message.contains("Connection refused")
|| message.contains("Connection reset")
|| message.contains("Connection timed out")
|| message.contains("Read timed out")
|| message.contains("No route to host")
|| message.contains("Service Unavailable")
|| message.contains("circuit_breaking_exception");
}
/**
* 创建索引不存在异常
*/
public static ElasticsearchException indexNotFound(String index) {
ElasticsearchException ex = new ElasticsearchException(
"INDEX_NOT_FOUND", index, "索引不存在");
ex.errorCode = "INDEX_NOT_FOUND";
return ex;
}
/**
* 创建文档不存在异常
*/
public static ElasticsearchException documentNotFound(String index, String id) {
ElasticsearchException ex = new ElasticsearchException(
"DOCUMENT_NOT_FOUND", index,
String.format("文档 [%s] 不存在", id));
ex.errorCode = "DOCUMENT_NOT_FOUND";
return ex;
}
/**
* 创建参数校验异常
*/
public static ElasticsearchException invalidParameter(String message) {
ElasticsearchException ex = new ElasticsearchException(message);
ex.errorCode = "INVALID_PARAMETER";
return ex;
}
}
```
### **9.6 结果封装类**
```java
package com.example.elasticsearch.result;
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;
import java.util.Collections;
import java.util.List;
/**
* 分页结果封装
*/
@Data
@NoArgsConstructor
@AllArgsConstructor
public class PageResult {
/**
* 数据列表
*/
private List records;
/**
* 总记录数
*/
private Long total;
/**
* 当前页码
*/
private Integer pageNum;
/**
* 每页大小
*/
private Integer pageSize;
/**
* 总页数
*/
private Integer totalPages;
/**
* 是否有下一页
*/
private Boolean hasNext;
/**
* 是否有上一页
*/
private Boolean hasPrevious;
public PageResult(List records, Long total, Integer pageNum, Integer pageSize) {
this.records = records;
this.total = total;
this.pageNum = pageNum;
this.pageSize = pageSize;
this.totalPages = (int) Math.ceil((double) total / pageSize);
this.hasNext = pageNum < totalPages;
this.hasPrevious = pageNum > 1;
}
/**
* 空结果
*/
public static PageResult empty(Integer pageNum, Integer pageSize) {
return new PageResult<>(Collections.emptyList(), 0L, pageNum, pageSize);
}
}
```
```java
package com.example.elasticsearch.result;
import lombok.Data;
import java.util.List;
import java.util.Map;
/**
* 搜索结果封装(包含高亮、评分等信息)
*/
@Data
public class SearchResult {
/**
* 文档 ID
*/
private String documentId;
/**
* 数据对象
*/
private T source;
/**
* 评分
*/
private Float score;
/**
* 高亮字段
*/
private Map> highlight;
/**
* 排序值(用于深度分页)
*/
private Object[] sortValues;
public SearchResult(String documentId, T source) {
this.documentId = documentId;
this.source = source;
}
public SearchResult(String documentId, T source, Float score,
Map> highlight) {
this.documentId = documentId;
this.source = source;
this.score = score;
this.highlight = highlight;
}
}
```
```java
package com.example.elasticsearch.result;
import lombok.Data;
import java.util.List;
import java.util.Map;
/**
* 聚合结果封装
*/
@Data
public class AggregationResult {
/**
* 聚合名称
*/
private String name;
/**
* 桶数据(terms 聚合)
*/
private List buckets;
/**
* 数值(sum、avg、max、min 等)
*/
private Double value;
@Data
public static class BucketData {
private String key;
private Long docCount;
private Map subAggregations;
}
}
```
### **9.7 查询构建器**
```java
package com.example.elasticsearch.query;
import lombok.Data;
import org.elasticsearch.index.query.*;
import org.elasticsearch.search.aggregations.AggregationBuilder;
import org.elasticsearch.search.aggregations.AggregationBuilders;
import org.elasticsearch.search.builder.SearchSourceBuilder;
import org.elasticsearch.search.fetch.subphase.highlight.HighlightBuilder;
import org.elasticsearch.search.sort.SortBuilder;
import org.elasticsearch.search.sort.SortBuilders;
import org.elasticsearch.search.sort.SortOrder;
import org.springframework.util.CollectionUtils;
import org.springframework.util.StringUtils;
import java.util.*;
/**
* ES 查询构建器
* 使用链式调用构建复杂查询
*/
@Data
public class EsQueryBuilder {
/**
* 索引名称
*/
private String indexName;
/**
* Bool 查询条件
*/
private BoolQueryBuilder boolQuery;
/**
* 排序条件
*/
private List> sorts;
/**
* 高亮配置
*/
private HighlightBuilder highlightBuilder;
/**
* 聚合配置
*/
private List aggregations;
/**
* 分页参数
*/
private Integer from;
private Integer size;
/**
* 返回字段(包含)
*/
private String[] includes;
/**
* 排除字段
*/
private String[] excludes;
/**
* Search After(深度分页)
*/
private Object[] searchAfter;
/**
* 是否追踪总数
*/
private Boolean trackTotalHits = true;
private EsQueryBuilder() {
this.boolQuery = QueryBuilders.boolQuery();
this.sorts = new ArrayList<>();
this.aggregations = new ArrayList<>();
}
/**
* 创建构建器
*/
public static EsQueryBuilder builder(String indexName) {
EsQueryBuilder builder = new EsQueryBuilder();
builder.indexName = indexName;
return builder;
}
// ==================== Must 条件(必须匹配)====================
/**
* 精确匹配(term)
*/
public EsQueryBuilder term(String field, Object value) {
if (value != null) {
boolQuery.must(QueryBuilders.termQuery(field, value));
}
return this;
}
/**
* 多值匹配(terms)
*/
public EsQueryBuilder terms(String field, Collection> values) {
if (!CollectionUtils.isEmpty(values)) {
boolQuery.must(QueryBuilders.termsQuery(field, values));
}
return this;
}
/**
* 全文匹配(match)
*/
public EsQueryBuilder match(String field, Object value) {
if (value != null && StringUtils.hasText(value.toString())) {
boolQuery.must(QueryBuilders.matchQuery(field, value));
}
return this;
}
/**
* 短语匹配(match_phrase)
*/
public EsQueryBuilder matchPhrase(String field, Object value) {
if (value != null && StringUtils.hasText(value.toString())) {
boolQuery.must(QueryBuilders.matchPhraseQuery(field, value));
}
return this;
}
/**
* 多字段匹配(multi_match)
*/
public EsQueryBuilder multiMatch(Object value, String... fields) {
if (value != null && StringUtils.hasText(value.toString()) && fields.length > 0) {
boolQuery.must(QueryBuilders.multiMatchQuery(value, fields));
}
return this;
}
/**
* 前缀匹配(prefix)
*/
public EsQueryBuilder prefix(String field, String prefix) {
if (StringUtils.hasText(prefix)) {
boolQuery.must(QueryBuilders.prefixQuery(field, prefix));
}
return this;
}
/**
* 通配符匹配(wildcard)
*/
public EsQueryBuilder wildcard(String field, String pattern) {
if (StringUtils.hasText(pattern)) {
boolQuery.must(QueryBuilders.wildcardQuery(field, pattern));
}
return this;
}
/**
* 范围查询
*/
public EsQueryBuilder range(String field, Object gte, Object lte) {
RangeQueryBuilder rangeQuery = QueryBuilders.rangeQuery(field);
if (gte != null) {
rangeQuery.gte(gte);
}
if (lte != null) {
rangeQuery.lte(lte);
}
if (gte != null || lte != null) {
boolQuery.must(rangeQuery);
}
return this;
}
/**
* 范围查询(大于)
*/
public EsQueryBuilder gt(String field, Object value) {
if (value != null) {
boolQuery.must(QueryBuilders.rangeQuery(field).gt(value));
}
return this;
}
/**
* 范围查询(大于等于)
*/
public EsQueryBuilder gte(String field, Object value) {
if (value != null) {
boolQuery.must(QueryBuilders.rangeQuery(field).gte(value));
}
return this;
}
/**
* 范围查询(小于)
*/
public EsQueryBuilder lt(String field, Object value) {
if (value != null) {
boolQuery.must(QueryBuilders.rangeQuery(field).lt(value));
}
return this;
}
/**
* 范围查询(小于等于)
*/
public EsQueryBuilder lte(String field, Object value) {
if (value != null) {
boolQuery.must(QueryBuilders.rangeQuery(field).lte(value));
}
return this;
}
/**
* 存在字段查询
*/
public EsQueryBuilder exists(String field) {
boolQuery.must(QueryBuilders.existsQuery(field));
return this;
}
// ==================== Filter 条件(过滤,不计算评分)====================
/**
* Filter - 精确匹配
*/
public EsQueryBuilder filterTerm(String field, Object value) {
if (value != null) {
boolQuery.filter(QueryBuilders.termQuery(field, value));
}
return this;
}
/**
* Filter - 多值匹配
*/
public EsQueryBuilder filterTerms(String field, Collection> values) {
if (!CollectionUtils.isEmpty(values)) {
boolQuery.filter(QueryBuilders.termsQuery(field, values));
}
return this;
}
/**
* Filter - 范围查询
*/
public EsQueryBuilder filterRange(String field, Object gte, Object lte) {
RangeQueryBuilder rangeQuery = QueryBuilders.rangeQuery(field);
if (gte != null) {
rangeQuery.gte(gte);
}
if (lte != null) {
rangeQuery.lte(lte);
}
if (gte != null || lte != null) {
boolQuery.filter(rangeQuery);
}
return this;
}
// ==================== Should 条件(或条件)====================
/**
* Should - 至少匹配一个
*/
public EsQueryBuilder should(QueryBuilder... queries) {
for (QueryBuilder query : queries) {
boolQuery.should(query);
}
return this;
}
/**
* Should - 多字段或查询
*/
public EsQueryBuilder shouldMatch(String value, String... fields) {
if (StringUtils.hasText(value) && fields.length > 0) {
for (String field : fields) {
boolQuery.should(QueryBuilders.matchQuery(field, value));
}
boolQuery.minimumShouldMatch(1);
}
return this;
}
/**
* 设置最小 should 匹配数
*/
public EsQueryBuilder minimumShouldMatch(int count) {
boolQuery.minimumShouldMatch(count);
return this;
}
// ==================== MustNot 条件(必须不匹配)====================
/**
* MustNot - 精确匹配
*/
public EsQueryBuilder mustNotTerm(String field, Object value) {
if (value != null) {
boolQuery.mustNot(QueryBuilders.termQuery(field, value));
}
return this;
}
/**
* MustNot - 多值匹配
*/
public EsQueryBuilder mustNotTerms(String field, Collection> values) {
if (!CollectionUtils.isEmpty(values)) {
boolQuery.mustNot(QueryBuilders.termsQuery(field, values));
}
return this;
}
// ==================== 嵌套查询 ====================
/**
* 嵌套查询
*/
public EsQueryBuilder nested(String path, QueryBuilder query) {
boolQuery.must(QueryBuilders.nestedQuery(path, query,
org.apache.lucene.search.join.ScoreMode.Avg));
return this;
}
/**
* 添加自定义查询条件
*/
public EsQueryBuilder must(QueryBuilder query) {
boolQuery.must(query);
return this;
}
/**
* 添加自定义过滤条件
*/
public EsQueryBuilder filter(QueryBuilder query) {
boolQuery.filter(query);
return this;
}
// ==================== 排序 ====================
/**
* 添加排序
*/
public EsQueryBuilder sort(String field, SortOrder order) {
sorts.add(SortBuilders.fieldSort(field).order(order));
return this;
}
/**
* 按评分排序
*/
public EsQueryBuilder sortByScore(SortOrder order) {
sorts.add(SortBuilders.scoreSort().order(order));
return this;
}
/**
* 多字段排序
*/
public EsQueryBuilder sorts(Map sortMap) {
sortMap.forEach((field, order) ->
sorts.add(SortBuilders.fieldSort(field).order(order)));
return this;
}
// ==================== 分页 ====================
/**
* 设置分页
*/
public EsQueryBuilder page(int pageNum, int pageSize) {
this.from = (pageNum - 1) * pageSize;
this.size = pageSize;
return this;
}
/**
* 设置起始位置和大小
*/
public EsQueryBuilder fromSize(int from, int size) {
this.from = from;
this.size = size;
return this;
}
/**
* 深度分页(Search After)
*/
public EsQueryBuilder searchAfter(Object[] values) {
this.searchAfter = values;
return this;
}
// ==================== 高亮 ====================
/**
* 添加高亮字段
*/
public EsQueryBuilder highlight(String... fields) {
if (fields.length > 0) {
highlightBuilder = new HighlightBuilder();
for (String field : fields) {
highlightBuilder.field(field);
}
highlightBuilder.preTags("");
highlightBuilder.postTags("");
}
return this;
}
/**
* 自定义高亮配置
*/
public EsQueryBuilder highlight(String preTag, String postTag, String... fields) {
if (fields.length > 0) {
highlightBuilder = new HighlightBuilder();
for (String field : fields) {
highlightBuilder.field(field);
}
highlightBuilder.preTags(preTag);
highlightBuilder.postTags(postTag);
}
return this;
}
// ==================== 聚合 ====================
/**
* Terms 聚合
*/
public EsQueryBuilder termsAggregation(String name, String field, int size) {
aggregations.add(AggregationBuilders.terms(name).field(field).size(size));
return this;
}
/**
* Sum 聚合
*/
public EsQueryBuilder sumAggregation(String name, String field) {
aggregations.add(AggregationBuilders.sum(name).field(field));
return this;
}
/**
* Avg 聚合
*/
public EsQueryBuilder avgAggregation(String name, String field) {
aggregations.add(AggregationBuilders.avg(name).field(field));
return this;
}
/**
* Max 聚合
*/
public EsQueryBuilder maxAggregation(String name, String field) {
aggregations.add(AggregationBuilders.max(name).field(field));
return this;
}
/**
* Min 聚合
*/
public EsQueryBuilder minAggregation(String name, String field) {
aggregations.add(AggregationBuilders.min(name).field(field));
return this;
}
/**
* 日期直方图聚合
*/
public EsQueryBuilder dateHistogramAggregation(String name, String field,
String interval) {
aggregations.add(AggregationBuilders.dateHistogram(name)
.field(field)
.calendarInterval(new org.elasticsearch.search.aggregations.bucket
.histogram.DateHistogramInterval(interval)));
return this;
}
/**
* 添加自定义聚合
*/
public EsQueryBuilder aggregation(AggregationBuilder aggregation) {
aggregations.add(aggregation);
return this;
}
// ==================== 返回字段 ====================
/**
* 指定返回字段
*/
public EsQueryBuilder includes(String... fields) {
this.includes = fields;
return this;
}
/**
* 排除字段
*/
public EsQueryBuilder excludes(String... fields) {
this.excludes = fields;
return this;
}
/**
* 是否追踪总数
*/
public EsQueryBuilder trackTotalHits(boolean track) {
this.trackTotalHits = track;
return this;
}
// ==================== 构建 ====================
/**
* 构建 SearchSourceBuilder
*/
public SearchSourceBuilder build() {
SearchSourceBuilder sourceBuilder = new SearchSourceBuilder();
// 查询条件
sourceBuilder.query(boolQuery);
// 分页
if (from != null) {
sourceBuilder.from(from);
}
if (size != null) {
sourceBuilder.size(size);
}
// 排序
for (SortBuilder> sort : sorts) {
sourceBuilder.sort(sort);
}
// 高亮
if (highlightBuilder != null) {
sourceBuilder.highlighter(highlightBuilder);
}
// 聚合
for (AggregationBuilder aggregation : aggregations) {
sourceBuilder.aggregation(aggregation);
}
// 返回字段
if (includes != null || excludes != null) {
sourceBuilder.fetchSource(includes, excludes);
}
// Search After
if (searchAfter != null) {
sourceBuilder.searchAfter(searchAfter);
}
// 追踪总数
sourceBuilder.trackTotalHits(trackTotalHits);
return sourceBuilder;
}
}
```
### **9.8 重试工具类**
```java
package com.example.elasticsearch.util;
import com.example.elasticsearch.exception.ElasticsearchException;
import lombok.extern.slf4j.Slf4j;
import java.util.concurrent.Callable;
import java.util.function.Predicate;
/**
* 重试工具类
*/
@Slf4j
public class RetryUtil {
/**
* 执行带重试的操作
*
* @param callable 要执行的操作
* @param maxRetries 最大重试次数
* @param retryInterval 重试间隔(毫秒)
* @param retryOn 判断是否需要重试的条件
* @param operationName 操作名称(用于日志)
* @return 操作结果
*/
public static T executeWithRetry(
Callable callable,
int maxRetries,
long retryInterval,
Predicate retryOn,
String operationName) {
Exception lastException = null;
for (int attempt = 0; attempt <= maxRetries; attempt++) {
try {
return callable.call();
} catch (Exception e) {
lastException = e;
// 判断是否需要重试
if (attempt < maxRetries && retryOn.test(e)) {
log.warn("[{}] 操作失败,第 {}/{} 次重试,错误: {}",
operationName, attempt + 1, maxRetries, e.getMessage());
try {
Thread.sleep(retryInterval * (attempt + 1)); // 指数退避
} catch (InterruptedException ie) {
Thread.currentThread().interrupt();
throw new ElasticsearchException("重试被中断", ie);
}
} else {
break;
}
}
}
log.error("[{}] 操作失败,已达最大重试次数", operationName, lastException);
if (lastException instanceof ElasticsearchException) {
throw (ElasticsearchException) lastException;
}
throw new ElasticsearchException(operationName + " 操作失败", lastException);
}
/**
* 执行带重试的操作(无返回值)
*/
public static void executeWithRetry(
Runnable runnable,
int maxRetries,
long retryInterval,
Predicate retryOn,
String operationName) {
executeWithRetry(() -> {
runnable.run();
return null;
}, maxRetries, retryInterval, retryOn, operationName);
}
/**
* 默认的重试条件判断
*/
public static Predicate defaultRetryCondition() {
return e -> {
if (e instanceof ElasticsearchException) {
return ((ElasticsearchException) e).isRetryable();
}
String message = e.getMessage();
if (message == null) {
return false;
}
return message.contains("Connection")
|| message.contains("timed out")
|| message.contains("Unavailable");
};
}
}
```
### **9.9 参数校验工具类**
```java
package com.example.elasticsearch.util;
import com.example.elasticsearch.exception.ElasticsearchException;
import org.springframework.util.CollectionUtils;
import org.springframework.util.StringUtils;
import java.util.Collection;
/**
* 参数校验工具类
*/
public class ParamValidator {
private ParamValidator() {}
/**
* 校验索引名称
*/
public static void validateIndexName(String indexName) {
if (!StringUtils.hasText(indexName)) {
throw ElasticsearchException.invalidParameter("索引名称不能为空");
}
if (indexName.contains(" ")) {
throw ElasticsearchException.invalidParameter("索引名称不能包含空格");
}
if (!indexName.equals(indexName.toLowerCase())) {
throw ElasticsearchException.invalidParameter("索引名称必须小写");
}
}
/**
* 校验文档ID
*/
public static void validateDocumentId(String documentId) {
if (!StringUtils.hasText(documentId)) {
throw ElasticsearchException.invalidParameter("文档ID不能为空");
}
}
/**
* 校验分页参数
*/
public static void validatePageParam(int pageNum, int pageSize) {
if (pageNum < 1) {
throw ElasticsearchException.invalidParameter("页码必须大于0");
}
if (pageSize < 1 || pageSize > 10000) {
throw ElasticsearchException.invalidParameter("每页大小必须在1-10000之间");
}
// ES 默认限制 from + size <= 10000
if ((long) (pageNum - 1) * pageSize + pageSize > 10000) {
throw ElasticsearchException.invalidParameter(
"分页深度超出限制,请使用 searchAfter 方式");
}
}
/**
* 校验批量数据
*/
public static void validateBatchData(Collection> dataList) {
if (CollectionUtils.isEmpty(dataList)) {
throw ElasticsearchException.invalidParameter("批量数据不能为空");
}
}
/**
* 校验非空
*/
public static void notNull(Object object, String message) {
if (object == null) {
throw ElasticsearchException.invalidParameter(message);
}
}
/**
* 校验字符串非空
*/
public static void notBlank(String str, String message) {
if (!StringUtils.hasText(str)) {
throw ElasticsearchException.invalidParameter(message);
}
}
}
```
### **9.10 核心 Service 封装**
```java
package com.example.elasticsearch.core;
import com.alibaba.fastjson.JSON;
import com.example.elasticsearch.config.ElasticsearchProperties;
import com.example.elasticsearch.exception.ElasticsearchException;
import com.example.elasticsearch.query.EsQueryBuilder;
import com.example.elasticsearch.result.AggregationResult;
import com.example.elasticsearch.result.PageResult;
import com.example.elasticsearch.result.ScrollResult;
import com.example.elasticsearch.result.SearchResult;
import com.example.elasticsearch.util.ParamValidator;
import com.example.elasticsearch.util.RetryUtil;
import lombok.extern.slf4j.Slf4j;
import org.elasticsearch.action.search.*;
import org.elasticsearch.client.RequestOptions;
import org.elasticsearch.client.RestHighLevelClient;
import org.elasticsearch.common.text.Text;
import org.elasticsearch.common.unit.TimeValue;
import org.elasticsearch.index.query.QueryBuilder;
import org.elasticsearch.search.SearchHit;
import org.elasticsearch.search.SearchHits;
import org.elasticsearch.search.aggregations.Aggregation;
import org.elasticsearch.search.aggregations.bucket.terms.Terms;
import org.elasticsearch.search.aggregations.metrics.*;
import org.elasticsearch.search.builder.SearchSourceBuilder;
import org.elasticsearch.search.fetch.subphase.highlight.HighlightField;
import java.io.IOException;
import java.util.*;
import java.util.stream.Collectors;
/**
* Elasticsearch 核心搜索服务
*
* 特性:
* - 完善的参数校验
* - 自动重试机制
* - 慢查询日志
* - 滚动查询支持
* - 空值安全处理
*/
@Slf4j
public class ElasticsearchService {
private final RestHighLevelClient client;
private final ElasticsearchProperties properties;
public ElasticsearchService(RestHighLevelClient client,
ElasticsearchProperties properties) {
this.client = client;
this.properties = properties;
}
// ==================== 基础查询 ====================
/**
* 简单查询 - 根据单个字段精确匹配
*
* @param indexName 索引名称
* @param field 字段名
* @param value 字段值
* @param clazz 返回类型
* @return 匹配的文档列表
*/
public List query(String indexName, String field, Object value,
Class clazz) {
ParamValidator.validateIndexName(indexName);
ParamValidator.notBlank(field, "查询字段不能为空");
EsQueryBuilder queryBuilder = EsQueryBuilder.builder(indexName)
.term(field, value);
return search(queryBuilder, clazz);
}
/**
* 多条件查询 - 根据多个字段匹配
*
* @param indexName 索引名称
* @param params 查询参数 (字段名 -> 字段值)
* @param clazz 返回类型
* @return 匹配的文档列表
*/
public List query(String indexName, Map params,
Class clazz) {
ParamValidator.validateIndexName(indexName);
EsQueryBuilder queryBuilder = EsQueryBuilder.builder(indexName);
if (params != null && !params.isEmpty()) {
params.forEach((field, value) -> {
if (value != null) {
if (value instanceof String && ((String) value).length() > 0) {
queryBuilder.match(field, value);
} else if (!(value instanceof String)) {
queryBuilder.filterTerm(field, value);
}
}
});
}
return search(queryBuilder, clazz);
}
/**
* 分页查询
*
* @param indexName 索引名称
* @param params 查询参数
* @param pageNum 页码(从1开始)
* @param pageSize 每页大小
* @param clazz 返回类型
* @return 分页结果
*/
public PageResult queryPage(String indexName, Map params,
int pageNum, int pageSize, Class clazz) {
ParamValidator.validateIndexName(indexName);
ParamValidator.validatePageParam(pageNum, pageSize);
EsQueryBuilder queryBuilder = EsQueryBuilder.builder(indexName)
.page(pageNum, pageSize);
if (params != null && !params.isEmpty()) {
params.forEach((field, value) -> {
if (value != null) {
if (value instanceof String && ((String) value).length() > 0) {
queryBuilder.match(field, value);
} else if (!(value instanceof String)) {
queryBuilder.filterTerm(field, value);
}
}
});
}
return searchPage(queryBuilder, pageNum, pageSize, clazz);
}
// ==================== 高级查询 ====================
/**
* 使用查询构建器执行查询
*
* @param queryBuilder 查询构建器
* @param clazz 返回类型
* @return 匹配的文档列表
*/
public List search(EsQueryBuilder queryBuilder, Class clazz) {
ParamValidator.validateIndexName(queryBuilder.getIndexName());
ParamValidator.notNull(clazz, "返回类型不能为空");
return RetryUtil.executeWithRetry(
() -> doSearch(queryBuilder, clazz),
properties.getRetryTimes(),
properties.getRetryInterval(),
RetryUtil.defaultRetryCondition(),
"SEARCH"
);
}
private List doSearch(EsQueryBuilder queryBuilder, Class clazz)
throws IOException {
long startTime = System.currentTimeMillis();
try {
SearchSourceBuilder sourceBuilder = queryBuilder.build();
SearchRequest searchRequest = new SearchRequest(queryBuilder.getIndexName());
searchRequest.source(sourceBuilder);
// 打印 DSL
if (properties.getPrintDsl()) {
log.info("ES 查询 DSL: {}", sourceBuilder.toString());
}
SearchResponse response = client.search(searchRequest, RequestOptions.DEFAULT);
// 检查响应状态
checkResponseStatus(response, queryBuilder.getIndexName());
return parseHits(response.getHits(), clazz);
} finally {
logSlowQuery(startTime, "SEARCH", queryBuilder.getIndexName());
}
}
/**
* 分页查询
*/
public PageResult searchPage(EsQueryBuilder queryBuilder,
int pageNum, int pageSize, Class clazz) {
ParamValidator.validateIndexName(queryBuilder.getIndexName());
ParamValidator.validatePageParam(pageNum, pageSize);
ParamValidator.notNull(clazz, "返回类型不能为空");
return RetryUtil.executeWithRetry(
() -> doSearchPage(queryBuilder, pageNum, pageSize, clazz),
properties.getRetryTimes(),
properties.getRetryInterval(),
RetryUtil.defaultRetryCondition(),
"SEARCH_PAGE"
);
}
private PageResult doSearchPage(EsQueryBuilder queryBuilder,
int pageNum, int pageSize,
Class clazz) throws IOException {
long startTime = System.currentTimeMillis();
try {
queryBuilder.page(pageNum, pageSize);
SearchSourceBuilder sourceBuilder = queryBuilder.build();
SearchRequest searchRequest = new SearchRequest(queryBuilder.getIndexName());
searchRequest.source(sourceBuilder);
if (properties.getPrintDsl()) {
log.info("ES 分页查询 DSL: {}", sourceBuilder.toString());
}
SearchResponse response = client.search(searchRequest, RequestOptions.DEFAULT);
checkResponseStatus(response, queryBuilder.getIndexName());
List records = parseHits(response.getHits(), clazz);
long total = getTotalHits(response.getHits());
return new PageResult<>(records, total, pageNum, pageSize);
} finally {
logSlowQuery(startTime, "SEARCH_PAGE", queryBuilder.getIndexName());
}
}
/**
* 带高亮的查询
*/
public List> searchWithHighlight(EsQueryBuilder queryBuilder,
Class clazz) {
ParamValidator.validateIndexName(queryBuilder.getIndexName());
return RetryUtil.executeWithRetry(
() -> doSearchWithHighlight(queryBuilder, clazz),
properties.getRetryTimes(),
properties.getRetryInterval(),
RetryUtil.defaultRetryCondition(),
"SEARCH_HIGHLIGHT"
);
}
private List> doSearchWithHighlight(
EsQueryBuilder queryBuilder, Class clazz) throws IOException {
long startTime = System.currentTimeMillis();
try {
SearchSourceBuilder sourceBuilder = queryBuilder.build();
SearchRequest searchRequest = new SearchRequest(queryBuilder.getIndexName());
searchRequest.source(sourceBuilder);
if (properties.getPrintDsl()) {
log.info("ES 高亮查询 DSL: {}", sourceBuilder.toString());
}
SearchResponse response = client.search(searchRequest, RequestOptions.DEFAULT);
checkResponseStatus(response, queryBuilder.getIndexName());
return parseHitsWithHighlight(response.getHits(), clazz);
} finally {
logSlowQuery(startTime, "SEARCH_HIGHLIGHT", queryBuilder.getIndexName());
}
}
/**
* 带高亮的分页查询
*/
public PageResult> searchPageWithHighlight(
EsQueryBuilder queryBuilder, int pageNum, int pageSize, Class clazz) {
ParamValidator.validateIndexName(queryBuilder.getIndexName());
ParamValidator.validatePageParam(pageNum, pageSize);
return RetryUtil.executeWithRetry(
() -> doSearchPageWithHighlight(queryBuilder, pageNum, pageSize, clazz),
properties.getRetryTimes(),
properties.getRetryInterval(),
RetryUtil.defaultRetryCondition(),
"SEARCH_PAGE_HIGHLIGHT"
);
}
private PageResult> doSearchPageWithHighlight(
EsQueryBuilder queryBuilder, int pageNum, int pageSize,
Class clazz) throws IOException {
long startTime = System.currentTimeMillis();
try {
queryBuilder.page(pageNum, pageSize);
SearchSourceBuilder sourceBuilder = queryBuilder.build();
SearchRequest searchRequest = new SearchRequest(queryBuilder.getIndexName());
searchRequest.source(sourceBuilder);
SearchResponse response = client.search(searchRequest, RequestOptions.DEFAULT);
checkResponseStatus(response, queryBuilder.getIndexName());
List> records = parseHitsWithHighlight(response.getHits(), clazz);
long total = getTotalHits(response.getHits());
return new PageResult<>(records, total, pageNum, pageSize);
} finally {
logSlowQuery(startTime, "SEARCH_PAGE_HIGHLIGHT", queryBuilder.getIndexName());
}
}
// ==================== 滚动查询(大数据量) ====================
/**
* 滚动查询 - 初始化
* 适用于导出大量数据的场景
*
* @param queryBuilder 查询构建器
* @param scrollTime 滚动上下文保持时间(分钟)
* @param size 每批大小
* @param clazz 返回类型
* @return 滚动结果(包含scrollId和首批数据)
*/
public ScrollResult scrollSearch(EsQueryBuilder queryBuilder,
int scrollTime, int size,
Class clazz) {
ParamValidator.validateIndexName(queryBuilder.getIndexName());
try {
queryBuilder.fromSize(0, size);
SearchSourceBuilder sourceBuilder = queryBuilder.build();
SearchRequest searchRequest = new SearchRequest(queryBuilder.getIndexName());
searchRequest.source(sourceBuilder);
searchRequest.scroll(TimeValue.timeValueMinutes(scrollTime));
if (properties.getPrintDsl()) {
log.info("ES 滚动查询初始化 DSL: {}", sourceBuilder.toString());
}
SearchResponse response = client.search(searchRequest, RequestOptions.DEFAULT);
checkResponseStatus(response, queryBuilder.getIndexName());
List records = parseHits(response.getHits(), clazz);
long total = getTotalHits(response.getHits());
String scrollId = response.getScrollId();
return new ScrollResult<>(scrollId, records, total, records.size() < size);
} catch (IOException e) {
throw new ElasticsearchException("SCROLL_SEARCH",
queryBuilder.getIndexName(), e);
}
}
/**
* 滚动查询 - 继续获取下一批
*
* @param scrollId 滚动ID
* @param scrollTime 滚动上下文保持时间(分钟)
* @param size 每批大小
* @param clazz 返回类型
* @return 滚动结果
*/
public ScrollResult scrollNext(String scrollId, int scrollTime,
int size, Class clazz) {
ParamValidator.notBlank(scrollId, "scrollId不能为空");
try {
SearchScrollRequest scrollRequest = new SearchScrollRequest(scrollId);
scrollRequest.scroll(TimeValue.timeValueMinutes(scrollTime));
SearchResponse response = client.scroll(scrollRequest, RequestOptions.DEFAULT);
List records = parseHits(response.getHits(), clazz);
String newScrollId = response.getScrollId();
return new ScrollResult<>(newScrollId, records,
getTotalHits(response.getHits()), records.size() < size);
} catch (IOException e) {
throw new ElasticsearchException("滚动查询失败", e);
}
}
/**
* 清除滚动上下文
*
* @param scrollIds 滚动ID列表
*/
public void clearScroll(String... scrollIds) {
if (scrollIds == null || scrollIds.length == 0) {
return;
}
try {
ClearScrollRequest clearScrollRequest = new ClearScrollRequest();
clearScrollRequest.scrollIds(Arrays.asList(scrollIds));
ClearScrollResponse response = client.clearScroll(
clearScrollRequest, RequestOptions.DEFAULT);
if (!response.isSucceeded()) {
log.warn("清除滚动上下文失败");
}
} catch (IOException e) {
log.warn("清除滚动上下文异常", e);
}
}
// ==================== 聚合查询 ====================
/**
* 聚合查询
*/
public Map searchAggregation(EsQueryBuilder queryBuilder) {
ParamValidator.validateIndexName(queryBuilder.getIndexName());
return RetryUtil.executeWithRetry(
() -> doSearchAggregation(queryBuilder),
properties.getRetryTimes(),
properties.getRetryInterval(),
RetryUtil.defaultRetryCondition(),
"SEARCH_AGGREGATION"
);
}
private Map doSearchAggregation(
EsQueryBuilder queryBuilder) throws IOException {
long startTime = System.currentTimeMillis();
try {
queryBuilder.fromSize(0, 0);
SearchSourceBuilder sourceBuilder = queryBuilder.build();
SearchRequest searchRequest = new SearchRequest(queryBuilder.getIndexName());
searchRequest.source(sourceBuilder);
if (properties.getPrintDsl()) {
log.info("ES 聚合查询 DSL: {}", sourceBuilder.toString());
}
SearchResponse response = client.search(searchRequest, RequestOptions.DEFAULT);
checkResponseStatus(response, queryBuilder.getIndexName());
return parseAggregations(response.getAggregations());
} finally {
logSlowQuery(startTime, "SEARCH_AGGREGATION", queryBuilder.getIndexName());
}
}
// ==================== Search After(深度分页) ====================
/**
* Search After 查询
* 适用于深度分页场景,避免 from+size 的 10000 限制
*
* @param queryBuilder 查询构建器
* @param searchAfterValues 上一页最后一条记录的排序值
* @param size 每页大小
* @param clazz 返回类型
* @return 搜索结果列表
*/
public List> searchAfter(EsQueryBuilder queryBuilder,
Object[] searchAfterValues,
int size, Class clazz) {
ParamValidator.validateIndexName(queryBuilder.getIndexName());
return RetryUtil.executeWithRetry(
() -> doSearchAfter(queryBuilder, searchAfterValues, size, clazz),
properties.getRetryTimes(),
properties.getRetryInterval(),
RetryUtil.defaultRetryCondition(),
"SEARCH_AFTER"
);
}
private List> doSearchAfter(
EsQueryBuilder queryBuilder, Object[] searchAfterValues,
int size, Class clazz) throws IOException {
long startTime = System.currentTimeMillis();
try {
queryBuilder.fromSize(0, size);
if (searchAfterValues != null && searchAfterValues.length > 0) {
queryBuilder.searchAfter(searchAfterValues);
}
SearchSourceBuilder sourceBuilder = queryBuilder.build();
SearchRequest searchRequest = new SearchRequest(queryBuilder.getIndexName());
searchRequest.source(sourceBuilder);
if (properties.getPrintDsl()) {
log.info("ES Search After DSL: {}", sourceBuilder.toString());
}
SearchResponse response = client.search(searchRequest, RequestOptions.DEFAULT);
checkResponseStatus(response, queryBuilder.getIndexName());
return parseHitsWithHighlight(response.getHits(), clazz);
} finally {
logSlowQuery(startTime, "SEARCH_AFTER", queryBuilder.getIndexName());
}
}
// ==================== 统计查询 ====================
/**
* 统计符合条件的文档数量
*/
public long count(EsQueryBuilder queryBuilder) {
ParamValidator.validateIndexName(queryBuilder.getIndexName());
return RetryUtil.executeWithRetry(
() -> doCount(queryBuilder),
properties.getRetryTimes(),
properties.getRetryInterval(),
RetryUtil.defaultRetryCondition(),
"COUNT"
);
}
private long doCount(EsQueryBuilder queryBuilder) throws IOException {
queryBuilder.fromSize(0, 0);
SearchSourceBuilder sourceBuilder = queryBuilder.build();
SearchRequest searchRequest = new SearchRequest(queryBuilder.getIndexName());
searchRequest.source(sourceBuilder);
SearchResponse response = client.search(searchRequest, RequestOptions.DEFAULT);
checkResponseStatus(response, queryBuilder.getIndexName());
return getTotalHits(response.getHits());
}
/**
* 判断是否存在符合条件的文档
*/
public boolean exists(EsQueryBuilder queryBuilder) {
return count(queryBuilder) > 0;
}
// ==================== 原生查询 ====================
/**
* 执行原生 QueryBuilder 查询
* 适用于封装方法无法满足的复杂场景
*/
public List searchByQueryBuilder(String indexName,
QueryBuilder queryBuilder,
int size, Class clazz) {
ParamValidator.validateIndexName(indexName);
ParamValidator.notNull(queryBuilder, "查询条件不能为空");
try {
SearchSourceBuilder sourceBuilder = new SearchSourceBuilder();
sourceBuilder.query(queryBuilder);
sourceBuilder.size(size);
sourceBuilder.trackTotalHits(true);
SearchRequest searchRequest = new SearchRequest(indexName);
searchRequest.source(sourceBuilder);
SearchResponse response = client.search(searchRequest, RequestOptions.DEFAULT);
checkResponseStatus(response, indexName);
return parseHits(response.getHits(), clazz);
} catch (IOException e) {
throw new ElasticsearchException("SEARCH_BY_QUERY_BUILDER", indexName, e);
}
}
// ==================== 内部工具方法 ====================
/**
* 解析命中结果
*/
private List parseHits(SearchHits hits, Class clazz) {
if (hits == null || hits.getHits() == null) {
return Collections.emptyList();
}
List result = new ArrayList<>();
for (SearchHit hit : hits.getHits()) {
try {
String source = hit.getSourceAsString();
if (source != null && !source.isEmpty()) {
T obj = JSON.parseObject(source, clazz);
result.add(obj);
}
} catch (Exception e) {
log.warn("解析文档失败, id: {}, error: {}", hit.getId(), e.getMessage());
}
}
return result;
}
/**
* 解析带高亮的命中结果
*/
private List> parseHitsWithHighlight(SearchHits hits,
Class clazz) {
if (hits == null || hits.getHits() == null) {
return Collections.emptyList();
}
List> result = new ArrayList<>();
for (SearchHit hit : hits.getHits()) {
try {
String source = hit.getSourceAsString();
if (source == null || source.isEmpty()) {
continue;
}
T obj = JSON.parseObject(source, clazz);
// 解析高亮
Map> highlightMap = new HashMap<>();
Map highlightFields = hit.getHighlightFields();
if (highlightFields != null && !highlightFields.isEmpty()) {
highlightFields.forEach((field, highlightField) -> {
if (highlightField.getFragments() != null) {
List fragments = Arrays.stream(highlightField.getFragments())
.map(Text::string)
.collect(Collectors.toList());
highlightMap.put(field, fragments);
}
});
}
SearchResult searchResult = new SearchResult<>(
hit.getId(), obj, hit.getScore(), highlightMap);
searchResult.setSortValues(hit.getSortValues());
result.add(searchResult);
} catch (Exception e) {
log.warn("解析文档失败, id: {}, error: {}", hit.getId(), e.getMessage());
}
}
return result;
}
/**
* 解析聚合结果
*/
private Map parseAggregations(
org.elasticsearch.search.aggregations.Aggregations aggregations) {
Map result = new HashMap<>();
if (aggregations == null) {
return result;
}
for (Aggregation aggregation : aggregations) {
try {
AggregationResult aggResult = new AggregationResult();
aggResult.setName(aggregation.getName());
if (aggregation instanceof Terms) {
Terms terms = (Terms) aggregation;
List buckets = new ArrayList<>();
for (Terms.Bucket bucket : terms.getBuckets()) {
AggregationResult.BucketData bucketData =
new AggregationResult.BucketData();
bucketData.setKey(bucket.getKeyAsString());
bucketData.setDocCount(bucket.getDocCount());
buckets.add(bucketData);
}
aggResult.setBuckets(buckets);
} else if (aggregation instanceof Sum) {
aggResult.setValue(((Sum) aggregation).getValue());
} else if (aggregation instanceof Avg) {
aggResult.setValue(((Avg) aggregation).getValue());
} else if (aggregation instanceof Max) {
aggResult.setValue(((Max) aggregation).getValue());
} else if (aggregation instanceof Min) {
aggResult.setValue(((Min) aggregation).getValue());
} else if (aggregation instanceof ValueCount) {
aggResult.setValue((double) ((ValueCount) aggregation).getValue());
} else if (aggregation instanceof Cardinality) {
aggResult.setValue((double) ((Cardinality) aggregation).getValue());
}
result.put(aggregation.getName(), aggResult);
} catch (Exception e) {
log.warn("解析聚合结果失败, name: {}, error: {}",
aggregation.getName(), e.getMessage());
}
}
return result;
}
/**
* 获取总命中数(兼容不同版本)
*/
private long getTotalHits(SearchHits hits) {
if (hits == null || hits.getTotalHits() == null) {
return 0L;
}
return hits.getTotalHits().value;
}
/**
* 检查响应状态
*/
private void checkResponseStatus(SearchResponse response, String indexName) {
if (response == null) {
throw new ElasticsearchException("SEARCH", indexName, "响应为空");
}
if (response.isTimedOut()) {
log.warn("ES 查询超时, index: {}", indexName);
}
if (response.getShardFailures() != null && response.getShardFailures().length > 0) {
log.warn("ES 查询部分分片失败, index: {}, failures: {}",
indexName, response.getShardFailures().length);
}
}
/**
* 记录慢查询日志
*/
private void logSlowQuery(long startTime, String operation, String indexName) {
long elapsed = System.currentTimeMillis() - startTime;
if (elapsed >= properties.getSlowQueryThreshold()) {
log.warn("[慢查询] 操作: {}, 索引: {}, 耗时: {}ms",
operation, indexName, elapsed);
} else {
log.debug("ES 操作完成, 操作: {}, 索引: {}, 耗时: {}ms",
operation, indexName, elapsed);
}
}
}
```
### **9.11 文档操作服务**
```java
package com.example.elasticsearch.core;
import com.alibaba.fastjson.JSON;
import com.example.elasticsearch.config.ElasticsearchProperties;
import com.example.elasticsearch.exception.ElasticsearchException;
import com.example.elasticsearch.util.ParamValidator;
import com.example.elasticsearch.util.RetryUtil;
import lombok.extern.slf4j.Slf4j;
import org.elasticsearch.action.DocWriteResponse;
import org.elasticsearch.action.bulk.BulkItemResponse;
import org.elasticsearch.action.bulk.BulkRequest;
import org.elasticsearch.action.bulk.BulkResponse;
import org.elasticsearch.action.delete.DeleteRequest;
import org.elasticsearch.action.delete.DeleteResponse;
import org.elasticsearch.action.get.*;
import org.elasticsearch.action.index.IndexRequest;
import org.elasticsearch.action.index.IndexResponse;
import org.elasticsearch.action.support.WriteRequest;
import org.elasticsearch.action.update.UpdateRequest;
import org.elasticsearch.action.update.UpdateResponse;
import org.elasticsearch.client.RequestOptions;
import org.elasticsearch.client.RestHighLevelClient;
import org.elasticsearch.common.xcontent.XContentType;
import org.elasticsearch.index.query.QueryBuilder;
import org.elasticsearch.index.reindex.BulkByScrollResponse;
import org.elasticsearch.index.reindex.DeleteByQueryRequest;
import org.elasticsearch.index.reindex.UpdateByQueryRequest;
import org.elasticsearch.script.Script;
import org.elasticsearch.script.ScriptType;
import org.springframework.util.CollectionUtils;
import org.springframework.util.StringUtils;
import java.io.IOException;
import java.util.*;
/**
* Elasticsearch 文档操作服务(增强版)
*
* 特性:
* - 完善的参数校验
* - 批量操作分批处理
* - 自动重试机制
* - 详细的操作日志
*/
@Slf4j
public class ElasticsearchDocumentService {
private final RestHighLevelClient client;
private final ElasticsearchProperties properties;
public ElasticsearchDocumentService(RestHighLevelClient client,
ElasticsearchProperties properties) {
this.client = client;
this.properties = properties;
}
// ==================== 新增操作 ====================
/**
* 添加文档(自动生成ID)
*/
public String add(String indexName, Object data) {
return add(indexName, null, data);
}
/**
* 添加文档(指定ID)
*/
public String add(String indexName, String id, Object data) {
ParamValidator.validateIndexName(indexName);
ParamValidator.notNull(data, "文档数据不能为空");
return RetryUtil.executeWithRetry(
() -> doAdd(indexName, id, data),
properties.getRetryTimes(),
properties.getRetryInterval(),
RetryUtil.defaultRetryCondition(),
"ADD_DOCUMENT"
);
}
private String doAdd(String indexName, String id, Object data) throws IOException {
IndexRequest request = new IndexRequest(indexName);
if (StringUtils.hasText(id)) {
request.id(id);
}
String jsonData = JSON.toJSONString(data);
request.source(jsonData, XContentType.JSON);
// 兼容 ES 6.x
if (properties.getEnableType()) {
request.type(properties.getDefaultType());
}
// 设置刷新策略
setRefreshPolicy(request);
IndexResponse response = client.index(request, RequestOptions.DEFAULT);
log.info("添加文档成功, index: {}, id: {}", indexName, response.getId());
return response.getId();
}
/**
* 批量添加文档
* 自动分批处理,避免大批量数据导致的内存溢出
*/
public BulkResult batchAdd(String indexName, List> dataList) {
ParamValidator.validateIndexName(indexName);
if (CollectionUtils.isEmpty(dataList)) {
return BulkResult.empty();
}
int batchSize = properties.getBulkBatchSize();
int totalSize = dataList.size();
int batchCount = (int) Math.ceil((double) totalSize / batchSize);
BulkResult.Builder resultBuilder = BulkResult.builder();
for (int i = 0; i < batchCount; i++) {
int fromIndex = i * batchSize;
int toIndex = Math.min(fromIndex + batchSize, totalSize);
List> batchData = dataList.subList(fromIndex, toIndex);
try {
BulkResult batchResult = doBatchAdd(indexName, batchData);
resultBuilder.merge(batchResult);
log.debug("批量添加进度: {}/{}, 成功: {}, 失败: {}",
toIndex, totalSize,
batchResult.getSuccessCount(),
batchResult.getFailureCount());
} catch (Exception e) {
log.error("批量添加第 {} 批失败", i + 1, e);
resultBuilder.addFailure(batchData.size(), e.getMessage());
}
}
BulkResult result = resultBuilder.build();
log.info("批量添加完成, index: {}, 总数: {}, 成功: {}, 失败: {}",
indexName, totalSize, result.getSuccessCount(), result.getFailureCount());
return result;
}
private BulkResult doBatchAdd(String indexName, List> dataList) throws IOException {
BulkRequest bulkRequest = new BulkRequest();
for (Object data : dataList) {
IndexRequest request = new IndexRequest(indexName);
request.source(JSON.toJSONString(data), XContentType.JSON);
if (properties.getEnableType()) {
request.type(properties.getDefaultType());
}
bulkRequest.add(request);
}
// 设置刷新策略
setRefreshPolicy(bulkRequest);
BulkResponse response = client.bulk(bulkRequest, RequestOptions.DEFAULT);
return parseBulkResponse(response);
}
/**
* 批量添加文档(带ID)
*/
public BulkResult batchAddWithId(String indexName, Map dataMap) {
ParamValidator.validateIndexName(indexName);
if (CollectionUtils.isEmpty(dataMap)) {
return BulkResult.empty();
}
int batchSize = properties.getBulkBatchSize();
List> entries = new ArrayList<>(dataMap.entrySet());
int totalSize = entries.size();
int batchCount = (int) Math.ceil((double) totalSize / batchSize);
BulkResult.Builder resultBuilder = BulkResult.builder();
for (int i = 0; i < batchCount; i++) {
int fromIndex = i * batchSize;
int toIndex = Math.min(fromIndex + batchSize, totalSize);
List> batchEntries = entries.subList(fromIndex, toIndex);
try {
BulkResult batchResult = doBatchAddWithId(indexName, batchEntries);
resultBuilder.merge(batchResult);
} catch (Exception e) {
log.error("批量添加第 {} 批失败", i + 1, e);
resultBuilder.addFailure(batchEntries.size(), e.getMessage());
}
}
BulkResult result = resultBuilder.build();
log.info("批量添加完成, index: {}, 总数: {}, 成功: {}, 失败: {}",
indexName, totalSize, result.getSuccessCount(), result.getFailureCount());
return result;
}
private BulkResult doBatchAddWithId(String indexName,
List> entries)
throws IOException {
BulkRequest bulkRequest = new BulkRequest();
for (Map.Entry entry : entries) {
IndexRequest request = new IndexRequest(indexName);
request.id(entry.getKey());
request.source(JSON.toJSONString(entry.getValue()), XContentType.JSON);
if (properties.getEnableType()) {
request.type(properties.getDefaultType());
}
bulkRequest.add(request);
}
setRefreshPolicy(bulkRequest);
BulkResponse response = client.bulk(bulkRequest, RequestOptions.DEFAULT);
return parseBulkResponse(response);
}
// ==================== 查询操作 ====================
/**
* 根据ID查询
*/
public Optional getById(String indexName, String id, Class clazz) {
ParamValidator.validateIndexName(indexName);
ParamValidator.validateDocumentId(id);
return RetryUtil.executeWithRetry(
() -> doGetById(indexName, id, clazz),
properties.getRetryTimes(),
properties.getRetryInterval(),
RetryUtil.defaultRetryCondition(),
"GET_BY_ID"
);
}
private Optional doGetById(String indexName, String id, Class clazz)
throws IOException {
GetRequest request = new GetRequest(indexName, id);
GetResponse response = client.get(request, RequestOptions.DEFAULT);
if (!response.isExists()) {
return Optional.empty();
}
String source = response.getSourceAsString();
if (source == null || source.isEmpty()) {
return Optional.empty();
}
T obj = JSON.parseObject(source, clazz);
return Optional.of(obj);
}
/**
* 批量根据ID查询
*/
public Map getByIds(String indexName, List ids, Class clazz) {
ParamValidator.validateIndexName(indexName);
if (CollectionUtils.isEmpty(ids)) {
return Collections.emptyMap();
}
try {
MultiGetRequest request = new MultiGetRequest();
for (String id : ids) {
request.add(new MultiGetRequest.Item(indexName, id));
}
MultiGetResponse response = client.mget(request, RequestOptions.DEFAULT);
Map result = new HashMap<>();
for (MultiGetItemResponse itemResponse : response.getResponses()) {
if (!itemResponse.isFailed() && itemResponse.getResponse().isExists()) {
String source = itemResponse.getResponse().getSourceAsString();
if (source != null && !source.isEmpty()) {
T obj = JSON.parseObject(source, clazz);
result.put(itemResponse.getId(), obj);
}
}
}
return result;
} catch (IOException e) {
throw new ElasticsearchException("GET_BY_IDS", indexName, e);
}
}
/**
* 判断文档是否存在
*/
public boolean exists(String indexName, String id) {
ParamValidator.validateIndexName(indexName);
ParamValidator.validateDocumentId(id);
try {
GetRequest request = new GetRequest(indexName, id);
request.fetchSourceContext(
org.elasticsearch.search.fetch.subphase.FetchSourceContext.DO_NOT_FETCH_SOURCE);
GetResponse response = client.get(request, RequestOptions.DEFAULT);
return response.isExists();
} catch (IOException e) {
throw new ElasticsearchException("EXISTS", indexName, e);
}
}
// ==================== 更新操作 ====================
/**
* 更新文档(全量更新,不存在则新增)
*/
public boolean upsert(String indexName, String id, Object data) {
ParamValidator.validateIndexName(indexName);
ParamValidator.validateDocumentId(id);
ParamValidator.notNull(data, "更新数据不能为空");
return RetryUtil.executeWithRetry(
() -> doUpsert(indexName, id, data),
properties.getRetryTimes(),
properties.getRetryInterval(),
RetryUtil.defaultRetryCondition(),
"UPSERT"
);
}
private boolean doUpsert(String indexName, String id, Object data) throws IOException {
UpdateRequest request = new UpdateRequest(indexName, id);
request.doc(JSON.toJSONString(data), XContentType.JSON);
request.docAsUpsert(true);
setRefreshPolicy(request);
UpdateResponse response = client.update(request, RequestOptions.DEFAULT);
log.info("更新文档成功, index: {}, id: {}, result: {}",
indexName, id, response.getResult());
return response.getResult() != DocWriteResponse.Result.NOOP;
}
/**
* 部分更新文档
*/
public boolean partialUpdate(String indexName, String id, Map fields) {
ParamValidator.validateIndexName(indexName);
ParamValidator.validateDocumentId(id);
if (CollectionUtils.isEmpty(fields)) {
return false;
}
return RetryUtil.executeWithRetry(
() -> doPartialUpdate(indexName, id, fields),
properties.getRetryTimes(),
properties.getRetryInterval(),
RetryUtil.defaultRetryCondition(),
"PARTIAL_UPDATE"
);
}
private boolean doPartialUpdate(String indexName, String id,
Map fields) throws IOException {
UpdateRequest request = new UpdateRequest(indexName, id);
request.doc(fields);
setRefreshPolicy(request);
UpdateResponse response = client.update(request, RequestOptions.DEFAULT);
log.info("部分更新文档成功, index: {}, id: {}", indexName, id);
return response.getResult() != DocWriteResponse.Result.NOOP;
}
/**
* 使用脚本更新
*/
public boolean updateByScript(String indexName, String id, String script,
Map params) {
ParamValidator.validateIndexName(indexName);
ParamValidator.validateDocumentId(id);
ParamValidator.notBlank(script, "脚本不能为空");
try {
UpdateRequest request = new UpdateRequest(indexName, id);
Script painlessScript = new Script(
ScriptType.INLINE,
"painless",
script,
params != null ? params : Collections.emptyMap()
);
request.script(painlessScript);
setRefreshPolicy(request);
UpdateResponse response = client.update(request, RequestOptions.DEFAULT);
log.info("脚本更新文档成功, index: {}, id: {}", indexName, id);
return response.getResult() != DocWriteResponse.Result.NOOP;
} catch (IOException e) {
throw new ElasticsearchException("UPDATE_BY_SCRIPT", indexName, e);
}
}
/**
* 根据条件批量更新
*/
public long updateByQuery(String indexName, QueryBuilder query,
String script, Map params) {
ParamValidator.validateIndexName(indexName);
ParamValidator.notNull(query, "查询条件不能为空");
ParamValidator.notBlank(script, "脚本不能为空");
try {
UpdateByQueryRequest request = new UpdateByQueryRequest(indexName);
request.setQuery(query);
request.setScript(new Script(
ScriptType.INLINE,
"painless",
script,
params != null ? params : Collections.emptyMap()
));
request.setRefresh(true);
BulkByScrollResponse response = client.updateByQuery(
request, RequestOptions.DEFAULT);
log.info("根据条件更新文档成功, index: {}, updated: {}",
indexName, response.getUpdated());
return response.getUpdated();
} catch (IOException e) {
throw new ElasticsearchException("UPDATE_BY_QUERY", indexName, e);
}
}
/**
* 批量更新
*/
public BulkResult batchUpdate(String indexName, Map dataMap) {
ParamValidator.validateIndexName(indexName);
if (CollectionUtils.isEmpty(dataMap)) {
return BulkResult.empty();
}
try {
BulkRequest bulkRequest = new BulkRequest();
dataMap.forEach((id, data) -> {
UpdateRequest request = new UpdateRequest(indexName, id);
request.doc(JSON.toJSONString(data), XContentType.JSON);
request.docAsUpsert(true);
bulkRequest.add(request);
});
setRefreshPolicy(bulkRequest);
BulkResponse response = client.bulk(bulkRequest, RequestOptions.DEFAULT);
BulkResult result = parseBulkResponse(response);
log.info("批量更新完成, index: {}, 成功: {}, 失败: {}",
indexName, result.getSuccessCount(), result.getFailureCount());
return result;
} catch (IOException e) {
throw new ElasticsearchException("BATCH_UPDATE", indexName, e);
}
}
// ==================== 删除操作 ====================
/**
* 删除文档
*/
public boolean delete(String indexName, String id) {
ParamValidator.validateIndexName(indexName);
ParamValidator.validateDocumentId(id);
return RetryUtil.executeWithRetry(
() -> doDelete(indexName, id),
properties.getRetryTimes(),
properties.getRetryInterval(),
RetryUtil.defaultRetryCondition(),
"DELETE"
);
}
private boolean doDelete(String indexName, String id) throws IOException {
DeleteRequest request = new DeleteRequest(indexName, id);
setRefreshPolicy(request);
DeleteResponse response = client.delete(request, RequestOptions.DEFAULT);
log.info("删除文档成功, index: {}, id: {}", indexName, id);
return response.getResult() == DocWriteResponse.Result.DELETED;
}
/**
* 批量删除
*/
public BulkResult batchDelete(String indexName, List ids) {
ParamValidator.validateIndexName(indexName);
if (CollectionUtils.isEmpty(ids)) {
return BulkResult.empty();
}
try {
BulkRequest bulkRequest = new BulkRequest();
for (String id : ids) {
DeleteRequest request = new DeleteRequest(indexName, id);
bulkRequest.add(request);
}
setRefreshPolicy(bulkRequest);
BulkResponse response = client.bulk(bulkRequest, RequestOptions.DEFAULT);
BulkResult result = parseBulkResponse(response);
log.info("批量删除完成, index: {}, 成功: {}, 失败: {}",
indexName, result.getSuccessCount(), result.getFailureCount());
return result;
} catch (IOException e) {
throw new ElasticsearchException("BATCH_DELETE", indexName, e);
}
}
/**
* 根据条件删除
*/
public long deleteByQuery(String indexName, QueryBuilder query) {
ParamValidator.validateIndexName(indexName);
ParamValidator.notNull(query, "查询条件不能为空");
try {
DeleteByQueryRequest request = new DeleteByQueryRequest(indexName);
request.setQuery(query);
request.setRefresh(true);
BulkByScrollResponse response = client.deleteByQuery(
request, RequestOptions.DEFAULT);
log.info("根据条件删除文档成功, index: {}, deleted: {}",
indexName, response.getDeleted());
return response.getDeleted();
} catch (IOException e) {
throw new ElasticsearchException("DELETE_BY_QUERY", indexName, e);
}
}
// ==================== 内部工具方法 ====================
/**
* 设置刷新策略
*/
private void setRefreshPolicy(IndexRequest request) {
String policy = properties.getBulkRefreshPolicy();
if ("immediate".equals(policy)) {
request.setRefreshPolicy(WriteRequest.RefreshPolicy.IMMEDIATE);
} else if ("wait_for".equals(policy)) {
request.setRefreshPolicy(WriteRequest.RefreshPolicy.WAIT_UNTIL);
}
}
private void setRefreshPolicy(UpdateRequest request) {
String policy = properties.getBulkRefreshPolicy();
if ("immediate".equals(policy)) {
request.setRefreshPolicy(WriteRequest.RefreshPolicy.IMMEDIATE);
} else if ("wait_for".equals(policy)) {
request.setRefreshPolicy(WriteRequest.RefreshPolicy.WAIT_UNTIL);
}
}
private void setRefreshPolicy(DeleteRequest request) {
String policy = properties.getBulkRefreshPolicy();
if ("immediate".equals(policy)) {
request.setRefreshPolicy(WriteRequest.RefreshPolicy.IMMEDIATE);
} else if ("wait_for".equals(policy)) {
request.setRefreshPolicy(WriteRequest.RefreshPolicy.WAIT_UNTIL);
}
}
private void setRefreshPolicy(BulkRequest request) {
String policy = properties.getBulkRefreshPolicy();
if ("immediate".equals(policy)) {
request.setRefreshPolicy(WriteRequest.RefreshPolicy.IMMEDIATE);
} else if ("wait_for".equals(policy)) {
request.setRefreshPolicy(WriteRequest.RefreshPolicy.WAIT_UNTIL);
}
}
/**
* 解析批量操作响应
*/
private BulkResult parseBulkResponse(BulkResponse response) {
BulkResult.Builder builder = BulkResult.builder();
for (BulkItemResponse itemResponse : response.getItems()) {
if (itemResponse.isFailed()) {
builder.addFailure(itemResponse.getId(),
itemResponse.getFailureMessage());
} else {
builder.addSuccess(itemResponse.getId());
}
}
return builder.build();
}
// ==================== 批量操作结果 ====================
/**
* 批量操作结果
*/
@lombok.Data
public static class BulkResult {
private int successCount;
private int failureCount;
private List successIds;
private List failures;
@lombok.Data
@lombok.AllArgsConstructor
public static class FailureItem {
private String id;
private String reason;
}
public static BulkResult empty() {
BulkResult result = new BulkResult();
result.successCount = 0;
result.failureCount = 0;
result.successIds = Collections.emptyList();
result.failures = Collections.emptyList();
return result;
}
public boolean hasFailures() {
return failureCount > 0;
}
public boolean isAllSuccess() {
return failureCount == 0;
}
public static Builder builder() {
return new Builder();
}
public static class Builder {
private List successIds = new ArrayList<>();
private List failures = new ArrayList<>();
public Builder addSuccess(String id) {
successIds.add(id);
return this;
}
public Builder addFailure(String id, String reason) {
failures.add(new FailureItem(id, reason));
return this;
}
public Builder addFailure(int count, String reason) {
for (int i = 0; i < count; i++) {
failures.add(new FailureItem(null, reason));
}
return this;
}
public Builder merge(BulkResult other) {
if (other.successIds != null) {
successIds.addAll(other.successIds);
}
if (other.failures != null) {
failures.addAll(other.failures);
}
return this;
}
public BulkResult build() {
BulkResult result = new BulkResult();
result.successIds = successIds;
result.failures = failures;
result.successCount = successIds.size();
result.failureCount = failures.size();
return result;
}
}
}
}
```
### **9.12 滚动查询结果类**
```java
package com.example.elasticsearch.result;
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;
import java.util.List;
/**
* 滚动查询结果
*/
@Data
@NoArgsConstructor
@AllArgsConstructor
public class ScrollResult {
/**
* 滚动ID(用于获取下一批数据)
*/
private String scrollId;
/**
* 当前批次数据
*/
private List records;
/**
* 总记录数
*/
private Long total;
/**
* 是否是最后一批(没有更多数据)
*/
private Boolean finished;
/**
* 判断是否还有更多数据
*/
public boolean hasMore() {
return !finished && records != null && !records.isEmpty();
}
}
```
### **9.13 索引操作服务**
```java
package com.example.elasticsearch.core;
import com.example.elasticsearch.config.ElasticsearchProperties;
import com.example.elasticsearch.exception.ElasticsearchException;
import lombok.extern.slf4j.Slf4j;
import org.elasticsearch.action.admin.indices.delete.DeleteIndexRequest;
import org.elasticsearch.action.support.master.AcknowledgedResponse;
import org.elasticsearch.client.RequestOptions;
import org.elasticsearch.client.RestHighLevelClient;
import org.elasticsearch.client.indices.*;
import org.elasticsearch.common.settings.Settings;
import org.elasticsearch.common.xcontent.XContentType;
import java.io.IOException;
import java.util.Map;
/**
* Elasticsearch 索引操作服务
*/
@Slf4j
public class ElasticsearchIndexService {
private final RestHighLevelClient client;
private final ElasticsearchProperties properties;
public ElasticsearchIndexService(RestHighLevelClient client,
ElasticsearchProperties properties) {
this.client = client;
this.properties = properties;
}
/**
* 创建索引
*/
public boolean createIndex(String indexName) {
return createIndex(indexName, null, null);
}
/**
* 创建索引(带 Mapping)
*/
public boolean createIndex(String indexName, String mapping) {
try {
CreateIndexRequest request = new CreateIndexRequest(indexName);
if (mapping != null) {
request.source(mapping, XContentType.JSON);
}
CreateIndexResponse response = client.indices()
.create(request, RequestOptions.DEFAULT);
log.info("创建索引成功, index: {}", indexName);
return response.isAcknowledged();
} catch (IOException e) {
log.error("创建索引失败", e);
throw new ElasticsearchException("CREATE_INDEX", indexName, e);
}
}
/**
* 创建索引(带 Settings 和 Mapping)
*/
public boolean createIndex(String indexName,
Map settings,
Map mapping) {
try {
CreateIndexRequest request = new CreateIndexRequest(indexName);
if (settings != null) {
request.settings(settings);
}
if (mapping != null) {
request.mapping(mapping);
}
CreateIndexResponse response = client.indices()
.create(request, RequestOptions.DEFAULT);
log.info("创建索引成功, index: {}", indexName);
return response.isAcknowledged();
} catch (IOException e) {
log.error("创建索引失败", e);
throw new ElasticsearchException("CREATE_INDEX", indexName, e);
}
}
/**
* 判断索引是否存在
*/
public boolean existsIndex(String indexName) {
try {
GetIndexRequest request = new GetIndexRequest(indexName);
return client.indices().exists(request, RequestOptions.DEFAULT);
} catch (IOException e) {
log.error("判断索引是否存在失败", e);
throw new ElasticsearchException("EXISTS_INDEX", indexName, e);
}
}
/**
* 删除索引
*/
public boolean deleteIndex(String indexName) {
try {
DeleteIndexRequest request = new DeleteIndexRequest(indexName);
AcknowledgedResponse response = client.indices()
.delete(request, RequestOptions.DEFAULT);
log.info("删除索引成功, index: {}", indexName);
return response.isAcknowledged();
} catch (IOException e) {
log.error("删除索引失败", e);
throw new ElasticsearchException("DELETE_INDEX", indexName, e);
}
}
/**
* 更新 Mapping
*/
public boolean updateMapping(String indexName, Map properties) {
try {
PutMappingRequest request = new PutMappingRequest(indexName);
request.source(Map.of("properties", properties));
AcknowledgedResponse response = client.indices()
.putMapping(request, RequestOptions.DEFAULT);
log.info("更新 Mapping 成功, index: {}", indexName);
return response.isAcknowledged();
} catch (IOException e) {
log.error("更新 Mapping 失败", e);
throw new ElasticsearchException("UPDATE_MAPPING", indexName, e);
}
}
/**
* 获取索引信息
*/
public GetIndexResponse getIndex(String indexName) {
try {
GetIndexRequest request = new GetIndexRequest(indexName);
return client.indices().get(request, RequestOptions.DEFAULT);
} catch (IOException e) {
log.error("获取索引信息失败", e);
throw new ElasticsearchException("GET_INDEX", indexName, e);
}
}
/**
* 刷新索引
*/
public void refreshIndex(String... indexNames) {
try {
org.elasticsearch.client.indices.RefreshRequest request =
new org.elasticsearch.client.indices.RefreshRequest(indexNames);
client.indices().refresh(request, RequestOptions.DEFAULT);
log.info("刷新索引成功, indices: {}", String.join(",", indexNames));
} catch (IOException e) {
log.error("刷新索引失败", e);
throw new ElasticsearchException("REFRESH_INDEX",
String.join(",", indexNames), e);
}
}
/**
* 创建索引别名
*/
public boolean createAlias(String indexName, String aliasName) {
try {
var request = new org.elasticsearch.client.indices.IndicesAliasesRequest();
var action = new org.elasticsearch.client.indices.IndicesAliasesRequest
.AliasActions(org.elasticsearch.client.indices.IndicesAliasesRequest
.AliasActions.Type.ADD)
.index(indexName)
.alias(aliasName);
request.addAliasAction(action);
var response = client.indices().updateAliases(request, RequestOptions.DEFAULT);
log.info("创建索引别名成功, index: {}, alias: {}", indexName, aliasName);
return response.isAcknowledged();
} catch (IOException e) {
log.error("创建索引别名失败", e);
throw new ElasticsearchException("CREATE_ALIAS", indexName, e);
}
}
}
```
### **9.14 架构图**
![[2-Learning/05-项目/08-企业级项目深读/02-damai_pro/05-业务深读/04-节目服务深度解析/assets/架构图-39a2cd40.jpg]]
1. ✅ **链式调用** - `EsQueryBuilder` 支持流畅的链式构建
2. ✅ **类型安全** - 泛型支持,自动转换结果类型
3. ✅ **功能完整** - 覆盖索引、文档、查询、聚合操作
4. ✅ **高亮支持** - 内置高亮配置和结果解析
5. ✅ **分页封装** - 统一的分页结果对象
6. ✅ **连接池** - 支持集群、认证、连接池配置
7. ✅ **版本兼容** - 支持 ES 6.x/7.x type 兼容
8. ✅ **异常处理** - 统一的异常封装
## **十、完整使用文档**
### **快速开始**
#### **添加依赖**
```text
com.example
elasticsearch-spring-boot-starter
1.0.0
org.springframework.boot
spring-boot-starter-data-elasticsearch
com.alibaba
fastjson
1.2.83
```
#### **配置文件**
```yaml
# application.yml
elasticsearch:
enabled: true
nodes:
- localhost:9200
username: elastic # 可选
password: your_password # 可选
scheme: http
# 连接配置
connect-timeout: 5000
socket-timeout: 30000
max-conn-total: 100
max-conn-per-route: 50
# 重试配置
retry-times: 3
retry-interval: 1000
# 批量操作配置
bulk-batch-size: 1000
bulk-refresh-policy: none # none, immediate, wait_for
# 日志配置
print-dsl: true # 开发环境建议开启
slow-query-threshold: 3000 # 慢查询阈值(ms)
# 版本兼容
enable-type: false # ES 7.x 设为 false
```
#### **创建实体类**
```java
package com.example.demo.entity;
import lombok.Data;
import java.util.Date;
/**
* 节目实体
*/
@Data
public class Program {
private Long id;
/**
* 节目标题
*/
private String title;
/**
* 演员
*/
private String actor;
/**
* 分类ID
*/
private Long categoryId;
/**
* 分类名称
*/
private String categoryName;
/**
* 地区ID
*/
private Long areaId;
/**
* 地区名称
*/
private String areaName;
/**
* 演出时间
*/
private Date showTime;
/**
* 最低价格
*/
private Double minPrice;
/**
* 最高价格
*/
private Double maxPrice;
/**
* 状态:1-上架, 0-下架
*/
private Integer status;
/**
* 封面图片
*/
private String coverImage;
/**
* 创建时间
*/
private Date createTime;
/**
* 更新时间
*/
private Date updateTime;
}
```
#### **创建搜索参数类**
```java
package com.example.demo.dto;
import lombok.Data;
import java.util.Date;
/**
* 节目搜索参数
*/
@Data
public class ProgramSearchParam {
/**
* 关键词(搜索标题和演员)
*/
private String keyword;
/**
* 分类ID
*/
private Long categoryId;
/**
* 地区ID
*/
private Long areaId;
/**
* 状态
*/
private Integer status;
/**
* 最低价格
*/
private Double minPrice;
/**
* 最高价格
*/
private Double maxPrice;
/**
* 开始时间
*/
private Date startTime;
/**
* 结束时间
*/
private Date endTime;
/**
* 页码
*/
private Integer pageNum = 1;
/**
* 每页大小
*/
private Integer pageSize = 10;
/**
* 排序字段
*/
private String sortField = "showTime";
/**
* 排序方式:asc, desc
*/
private String sortOrder = "asc";
}
```
### **完整业务示例**
```java
package com.example.demo.service;
import com.example.demo.dto.ProgramSearchParam;
import com.example.demo.entity.Program;
import com.example.elasticsearch.core.ElasticsearchDocumentService;
import com.example.elasticsearch.core.ElasticsearchDocumentService.BulkResult;
import com.example.elasticsearch.core.ElasticsearchIndexService;
import com.example.elasticsearch.core.ElasticsearchService;
import com.example.elasticsearch.query.EsQueryBuilder;
import com.example.elasticsearch.result.AggregationResult;
import com.example.elasticsearch.result.PageResult;
import com.example.elasticsearch.result.ScrollResult;
import com.example.elasticsearch.result.SearchResult;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.elasticsearch.search.sort.SortOrder;
import org.springframework.stereotype.Service;
import javax.annotation.PostConstruct;
import java.util.*;
import java.util.function.Consumer;
import java.util.stream.Collectors;
/**
* 节目搜索服务 - 完整示例
*/
@Slf4j
@Service
@RequiredArgsConstructor
public class ProgramSearchService {
private final ElasticsearchService esService;
private final ElasticsearchDocumentService documentService;
private final ElasticsearchIndexService indexService;
private static final String INDEX_NAME = "program";
// ==================== 索引管理 ====================
/**
* 初始化索引(应用启动时调用)
*/
@PostConstruct
public void initIndex() {
if (!indexService.existsIndex(INDEX_NAME)) {
String mapping = """
{
"mappings": {
"properties": {
"id": { "type": "long" },
"title": {
"type": "text",
"analyzer": "ik_max_word",
"search_analyzer": "ik_smart",
"fields": {
"keyword": { "type": "keyword" }
}
},
"actor": {
"type": "text",
"analyzer": "ik_max_word",
"fields": {
"keyword": { "type": "keyword" }
}
},
"categoryId": { "type": "long" },
"categoryName": { "type": "keyword" },
"areaId": { "type": "long" },
"areaName": { "type": "keyword" },
"showTime": { "type": "date" },
"minPrice": { "type": "double" },
"maxPrice": { "type": "double" },
"status": { "type": "integer" },
"coverImage": { "type": "keyword", "index": false },
"createTime": { "type": "date" },
"updateTime": { "type": "date" }
}
},
"settings": {
"number_of_shards": 3,
"number_of_replicas": 1,
"refresh_interval": "1s"
}
}
""";
indexService.createIndex(INDEX_NAME, mapping);
log.info("索引 [{}] 创建成功", INDEX_NAME);
}
}
/**
* 重建索引
*/
public void rebuildIndex() {
// 删除旧索引
if (indexService.existsIndex(INDEX_NAME)) {
indexService.deleteIndex(INDEX_NAME);
}
// 重新初始化
initIndex();
}
// ==================== 文档操作 ====================
/**
* 添加/更新节目
*/
public void saveProgram(Program program) {
program.setUpdateTime(new Date());
if (program.getCreateTime() == null) {
program.setCreateTime(new Date());
}
documentService.upsert(INDEX_NAME, String.valueOf(program.getId()), program);
}
/**
* 批量添加节目
*/
public BulkResult batchSavePrograms(List programs) {
Date now = new Date();
programs.forEach(p -> {
p.setUpdateTime(now);
if (p.getCreateTime() == null) {
p.setCreateTime(now);
}
});
return documentService.batchAdd(INDEX_NAME, programs);
}
/**
* 批量添加节目(带ID)
*/
public BulkResult batchSaveProgramsWithId(List programs) {
Date now = new Date();
Map dataMap = new HashMap<>();
programs.forEach(p -> {
p.setUpdateTime(now);
if (p.getCreateTime() == null) {
p.setCreateTime(now);
}
dataMap.put(String.valueOf(p.getId()), p);
});
return documentService.batchAddWithId(INDEX_NAME, dataMap);
}
/**
* 根据ID查询节目
*/
public Optional getProgramById(Long id) {
return documentService.getById(INDEX_NAME, String.valueOf(id), Program.class);
}
/**
* 批量查询节目
*/
public Map getProgramsByIds(List ids) {
List strIds = ids.stream()
.map(String::valueOf)
.collect(Collectors.toList());
Map resultMap = documentService.getByIds(
INDEX_NAME, strIds, Program.class);
// 转换 key 类型
return resultMap.entrySet().stream()
.collect(Collectors.toMap(
e -> Long.parseLong(e.getKey()),
Map.Entry::getValue
));
}
/**
* 删除节目
*/
public void deleteProgram(Long id) {
documentService.delete(INDEX_NAME, String.valueOf(id));
}
/**
* 批量删除节目
*/
public BulkResult batchDeletePrograms(List ids) {
List strIds = ids.stream()
.map(String::valueOf)
.collect(Collectors.toList());
return documentService.batchDelete(INDEX_NAME, strIds);
}
/**
* 更新节目状态
*/
public boolean updateProgramStatus(Long id, Integer status) {
Map fields = new HashMap<>();
fields.put("status", status);
fields.put("updateTime", new Date());
return documentService.partialUpdate(INDEX_NAME, String.valueOf(id), fields);
}
/**
* 批量更新节目状态(使用脚本)
*/
public long batchUpdateStatus(List ids, Integer status) {
EsQueryBuilder queryBuilder = EsQueryBuilder.builder(INDEX_NAME)
.filterTerms("id", ids);
String script = "ctx._source.status = params.status; " +
"ctx._source.updateTime = params.updateTime";
Map params = new HashMap<>();
params.put("status", status);
params.put("updateTime", System.currentTimeMillis());
return documentService.updateByQuery(
INDEX_NAME, queryBuilder.getBoolQuery(), script, params);
}
// ==================== 基础搜索 ====================
/**
* 根据分类查询节目列表
*/
public List findByCategory(Long categoryId) {
return esService.query(INDEX_NAME, "categoryId", categoryId, Program.class);
}
/**
* 根据多条件查询
*/
public List findByConditions(Long categoryId, Long areaId, Integer status) {
Map params = new HashMap<>();
params.put("categoryId", categoryId);
params.put("areaId", areaId);
params.put("status", status);
return esService.query(INDEX_NAME, params, Program.class);
}
/**
* 简单分页查询
*/
public PageResult findPage(int pageNum, int pageSize) {
return esService.queryPage(INDEX_NAME, null, pageNum, pageSize, Program.class);
}
// ==================== 高级搜索 ====================
/**
* 关键词搜索(标题 + 演员)
*/
public List searchByKeyword(String keyword) {
EsQueryBuilder queryBuilder = EsQueryBuilder.builder(INDEX_NAME)
.shouldMatch(keyword, "title", "actor")
.filterTerm("status", 1) // 只搜索上架的
.sort("showTime", SortOrder.ASC);
return esService.search(queryBuilder, Program.class);
}
/**
* 关键词搜索(带高亮)
*/
public List> searchWithHighlight(String keyword) {
EsQueryBuilder queryBuilder = EsQueryBuilder.builder(INDEX_NAME)
.shouldMatch(keyword, "title", "actor")
.filterTerm("status", 1)
.highlight("", "", "title", "actor")
.sort("showTime", SortOrder.ASC);
return esService.searchWithHighlight(queryBuilder, Program.class);
}
/**
* 复杂条件搜索 + 分页
*/
public PageResult search(ProgramSearchParam param) {
EsQueryBuilder queryBuilder = EsQueryBuilder.builder(INDEX_NAME)
// 关键词搜索(标题或演员)
.shouldMatch(param.getKeyword(), "title", "actor")
// 分类过滤
.filterTerm("categoryId", param.getCategoryId())
// 地区过滤
.filterTerm("areaId", param.getAreaId())
// 状态过滤
.filterTerm("status", param.getStatus())
// 价格范围
.filterRange("minPrice", param.getMinPrice(), param.getMaxPrice())
// 时间范围
.filterRange("showTime",
param.getStartTime() != null ? param.getStartTime().getTime() : null,
param.getEndTime() != null ? param.getEndTime().getTime() : null)
// 排序
.sort(param.getSortField(),
"desc".equalsIgnoreCase(param.getSortOrder())
? SortOrder.DESC : SortOrder.ASC);
return esService.searchPage(queryBuilder,
param.getPageNum(), param.getPageSize(), Program.class);
}
/**
* 复杂条件搜索 + 分页 + 高亮
*/
public PageResult> searchWithHighlight(ProgramSearchParam param) {
EsQueryBuilder queryBuilder = EsQueryBuilder.builder(INDEX_NAME)
.shouldMatch(param.getKeyword(), "title", "actor")
.filterTerm("categoryId", param.getCategoryId())
.filterTerm("areaId", param.getAreaId())
.filterTerm("status", param.getStatus())
.filterRange("minPrice", param.getMinPrice(), param.getMaxPrice())
.filterRange("showTime",
param.getStartTime() != null ? param.getStartTime().getTime() : null,
param.getEndTime() != null ? param.getEndTime().getTime() : null)
.highlight("title", "actor")
.sort(param.getSortField(),
"desc".equalsIgnoreCase(param.getSortOrder())
? SortOrder.DESC : SortOrder.ASC);
return esService.searchPageWithHighlight(queryBuilder,
param.getPageNum(), param.getPageSize(), Program.class);
}
// ==================== 深度分页(Search After) ====================
/**
* 使用 Search After 进行深度分页
* 适用于页码超过 1000 的场景
*
* @param param 搜索参数
* @param searchAfterValue 上一页最后一条记录的排序值
* @return 搜索结果(包含 sortValues 用于下一页查询)
*/
public List> searchAfter(ProgramSearchParam param,
Object[] searchAfterValue) {
EsQueryBuilder queryBuilder = EsQueryBuilder.builder(INDEX_NAME)
.shouldMatch(param.getKeyword(), "title", "actor")
.filterTerm("categoryId", param.getCategoryId())
.filterTerm("status", param.getStatus())
// 必须有排序字段,且最后一个字段应该是唯一的(如 _id)
.sort(param.getSortField(),
"desc".equalsIgnoreCase(param.getSortOrder())
? SortOrder.DESC : SortOrder.ASC)
.sort("id", SortOrder.ASC); // 使用 id 保证唯一性
return esService.searchAfter(queryBuilder, searchAfterValue,
param.getPageSize(), Program.class);
}
// ==================== 滚动查询(大数据导出) ====================
/**
* 滚动查询导出所有数据
* 适用于数据导出场景
*
* @param param 搜索参数
* @param consumer 数据消费者
*/
public void scrollExport(ProgramSearchParam param, Consumer> consumer) {
EsQueryBuilder queryBuilder = EsQueryBuilder.builder(INDEX_NAME)
.shouldMatch(param.getKeyword(), "title", "actor")
.filterTerm("categoryId", param.getCategoryId())
.filterTerm("status", param.getStatus())
.sort("id", SortOrder.ASC); // 滚动查询必须有排序
int batchSize = 1000;
int scrollTime = 5; // 滚动上下文保持时间(分钟)
// 初始化滚动查询
ScrollResult scrollResult = esService.scrollSearch(
queryBuilder, scrollTime, batchSize, Program.class);
try {
// 处理首批数据
if (!scrollResult.getRecords().isEmpty()) {
consumer.accept(scrollResult.getRecords());
}
// 继续获取后续数据
while (scrollResult.hasMore()) {
scrollResult = esService.scrollNext(
scrollResult.getScrollId(), scrollTime, batchSize, Program.class);
if (!scrollResult.getRecords().isEmpty()) {
consumer.accept(scrollResult.getRecords());
}
}
log.info("滚动导出完成, 总数: {}", scrollResult.getTotal());
} finally {
// 清除滚动上下文
esService.clearScroll(scrollResult.getScrollId());
}
}
// ==================== 聚合统计 ====================
/**
* 按分类统计节目数量
*/
public Map countByCategory() {
EsQueryBuilder queryBuilder = EsQueryBuilder.builder(INDEX_NAME)
.filterTerm("status", 1) // 只统计上架的
.termsAggregation("category_count", "categoryId", 100);
Map aggResult = esService.searchAggregation(queryBuilder);
AggregationResult categoryAgg = aggResult.get("category_count");
if (categoryAgg == null || categoryAgg.getBuckets() == null) {
return Collections.emptyMap();
}
return categoryAgg.getBuckets().stream()
.collect(Collectors.toMap(
AggregationResult.BucketData::getKey,
AggregationResult.BucketData::getDocCount
));
}
/**
* 统计价格区间
*/
public Map getPriceStats() {
EsQueryBuilder queryBuilder = EsQueryBuilder.builder(INDEX_NAME)
.filterTerm("status", 1)
.minAggregation("min_price", "minPrice")
.maxAggregation("max_price", "maxPrice")
.avgAggregation("avg_price", "minPrice");
Map aggResult = esService.searchAggregation(queryBuilder);
Map stats = new HashMap<>();
if (aggResult.containsKey("min_price")) {
stats.put("minPrice", aggResult.get("min_price").getValue());
}
if (aggResult.containsKey("max_price")) {
stats.put("maxPrice", aggResult.get("max_price").getValue());
}
if (aggResult.containsKey("avg_price")) {
stats.put("avgPrice", aggResult.get("avg_price").getValue());
}
return stats;
}
/**
* 综合统计
*/
public Map getStatistics() {
EsQueryBuilder queryBuilder = EsQueryBuilder.builder(INDEX_NAME)
.filterTerm("status", 1)
.termsAggregation("by_category", "categoryId", 50)
.termsAggregation("by_area", "areaId", 50)
.avgAggregation("avg_price", "minPrice")
.minAggregation("min_price", "minPrice")
.maxAggregation("max_price", "maxPrice");
Map aggResult = esService.searchAggregation(queryBuilder);
Map stats = new HashMap<>();
stats.put("aggregations", aggResult);
stats.put("total", esService.count(
EsQueryBuilder.builder(INDEX_NAME).filterTerm("status", 1)));
return stats;
}
// ==================== 工具方法 ====================
/**
* 统计符合条件的数量
*/
public long count(ProgramSearchParam param) {
EsQueryBuilder queryBuilder = EsQueryBuilder.builder(INDEX_NAME)
.shouldMatch(param.getKeyword(), "title", "actor")
.filterTerm("categoryId", param.getCategoryId())
.filterTerm("areaId", param.getAreaId())
.filterTerm("status", param.getStatus());
return esService.count(queryBuilder);
}
/**
* 判断是否存在符合条件的数据
*/
public boolean exists(ProgramSearchParam param) {
return count(param) > 0;
}
}
```
### **Controller 示例**
```java
package com.example.demo.controller;
import com.example.demo.dto.ProgramSearchParam;
import com.example.demo.entity.Program;
import com.example.demo.service.ProgramSearchService;
import com.example.elasticsearch.core.ElasticsearchDocumentService.BulkResult;
import com.example.elasticsearch.result.PageResult;
import com.example.elasticsearch.result.SearchResult;
import lombok.RequiredArgsConstructor;
import org.springframework.web.bind.annotation.*;
import java.util.List;
import java.util.Map;
import java.util.Optional;
/**
* 节目搜索 API
*/
@RestController
@RequestMapping("/api/programs")
@RequiredArgsConstructor
public class ProgramController {
private final ProgramSearchService searchService;
// ==================== 基础 CRUD ====================
/**
* 添加/更新节目
*/
@PostMapping
public String saveProgram(@RequestBody Program program) {
searchService.saveProgram(program);
return "success";
}
/**
* 批量添加节目
*/
@PostMapping("/batch")
public BulkResult batchSavePrograms(@RequestBody List programs) {
return searchService.batchSavePrograms(programs);
}
/**
* 根据ID查询
*/
@GetMapping("/{id}")
public Optional getProgram(@PathVariable Long id) {
return searchService.getProgramById(id);
}
/**
* 批量查询
*/
@PostMapping("/batch-get")
public Map batchGetPrograms(@RequestBody List ids) {
return searchService.getProgramsByIds(ids);
}
/**
* 删除节目
*/
@DeleteMapping("/{id}")
public String deleteProgram(@PathVariable Long id) {
searchService.deleteProgram(id);
return "success";
}
/**
* 更新状态
*/
@PutMapping("/{id}/status")
public boolean updateStatus(@PathVariable Long id, @RequestParam Integer status) {
return searchService.updateProgramStatus(id, status);
}
// ==================== 搜索 ====================
/**
* 关键词搜索
*/
@GetMapping("/search")
public List search(@RequestParam String keyword) {
return searchService.searchByKeyword(keyword);
}
/**
* 关键词搜索(带高亮)
*/
@GetMapping("/search/highlight")
public List> searchWithHighlight(@RequestParam String keyword) {
return searchService.searchWithHighlight(keyword);
}
/**
* 高级搜索(分页)
*/
@PostMapping("/search/page")
public PageResult searchPage(@RequestBody ProgramSearchParam param) {
return searchService.search(param);
}
/**
* 高级搜索(分页 + 高亮)
*/
@PostMapping("/search/page-highlight")
public PageResult> searchPageWithHighlight(
@RequestBody ProgramSearchParam param) {
return searchService.searchWithHighlight(param);
}
/**
* 深度分页(Search After)
*/
@PostMapping("/search/after")
public List> searchAfter(
@RequestBody ProgramSearchParam param,
@RequestParam(required = false) Long afterId,
@RequestParam(required = false) Long afterShowTime) {
Object[] searchAfter = null;
if (afterShowTime != null && afterId != null) {
searchAfter = new Object[]{afterShowTime, afterId};
}
return searchService.searchAfter(param, searchAfter);
}
// ==================== 统计 ====================
/**
* 按分类统计
*/
@GetMapping("/stats/category")
public Map countByCategory() {
return searchService.countByCategory();
}
/**
* 价格统计
*/
@GetMapping("/stats/price")
public Map getPriceStats() {
return searchService.getPriceStats();
}
/**
* 综合统计
*/
@GetMapping("/stats")
public Map getStatistics() {
return searchService.getStatistics();
}
// ==================== 索引管理 ====================
/**
* 重建索引
*/
@PostMapping("/index/rebuild")
public String rebuildIndex() {
searchService.rebuildIndex();
return "success";
}
}
```
### **最佳实践指南**
#### **索引设计**
```java
/**
* 索引设计最佳实践
*/
public class IndexDesignBestPractice {
/**
* 1. 合理设置分片数
* - 单个分片大小建议 10-50GB
* - 分片数 = 数据量 / 单分片大小
*/
public static final int SHARD_COUNT = 3;
/**
* 2. 副本数设置
* - 生产环境至少 1 个副本
* - 写入密集型可先设为 0,完成后再设为 1
*/
public static final int REPLICA_COUNT = 1;
/**
* 3. Mapping 设计原则
* - 明确字段类型,避免动态映射
* - text 用于全文搜索,keyword 用于精确匹配/聚合
* - 不需要搜索的字段设置 index: false
* - 大文本字段考虑使用 store: true
*/
public static String createOptimizedMapping() {
return """
{
"mappings": {
"properties": {
"id": { "type": "long" },
"title": {
"type": "text",
"analyzer": "ik_max_word",
"search_analyzer": "ik_smart",
"fields": {
"keyword": {
"type": "keyword",
"ignore_above": 256
}
}
},
"description": {
"type": "text",
"analyzer": "ik_max_word",
"index": true,
"store": true
},
"status": { "type": "keyword" },
"price": { "type": "scaled_float", "scaling_factor": 100 },
"createTime": { "type": "date", "format": "epoch_millis" },
"tags": { "type": "keyword" },
"coverImage": { "type": "keyword", "index": false }
}
},
"settings": {
"number_of_shards": 3,
"number_of_replicas": 1,
"refresh_interval": "1s",
"max_result_window": 10000,
"analysis": {
"analyzer": {
"ik_smart_pinyin": {
"type": "custom",
"tokenizer": "ik_smart",
"filter": ["lowercase", "pinyin_filter"]
}
}
}
}
}
""";
}
}
```
#### **查询优化**
```java
/**
* 查询优化最佳实践
*/
@Slf4j
@Service
public class QueryOptimizationService {
@Autowired
private ElasticsearchService esService;
private static final String INDEX_NAME = "program";
/**
* 1. 使用 filter 而非 must 进行过滤
* filter 不计算评分,性能更好,且结果可缓存
*/
public List optimizedFilterQuery(Long categoryId, Integer status) {
EsQueryBuilder queryBuilder = EsQueryBuilder.builder(INDEX_NAME)
// ✅ 正确:使用 filter 进行精确匹配
.filterTerm("categoryId", categoryId)
.filterTerm("status", status);
// ❌ 错误:使用 must 进行精确匹配会计算评分,浪费性能
// .term("categoryId", categoryId)
// .term("status", status)
return esService.search(queryBuilder, Program.class);
}
/**
* 2. 合理使用分页,避免深度分页
* from + size 超过 10000 会报错
*/
public PageResult safePageQuery(int pageNum, int pageSize) {
// 计算深度
int depth = (pageNum - 1) * pageSize + pageSize;
if (depth > 10000) {
// 深度分页使用 search_after
log.warn("分页深度超过10000,建议使用 searchAfter");
throw new IllegalArgumentException("分页深度超出限制,请使用 searchAfter");
}
EsQueryBuilder queryBuilder = EsQueryBuilder.builder(INDEX_NAME)
.page(pageNum, pageSize)
.sort("id", SortOrder.ASC);
return esService.searchPage(queryBuilder, pageNum, pageSize, Program.class);
}
/**
* 3. 只返回需要的字段,减少网络传输
*/
public List selectiveFieldsQuery(String keyword) {
EsQueryBuilder queryBuilder = EsQueryBuilder.builder(INDEX_NAME)
.match("title", keyword)
// 只返回需要的字段
.includes("id", "title", "minPrice", "showTime")
// 排除大字段
.excludes("description", "coverImage");
return esService.search(queryBuilder, Program.class);
}
/**
* 4. 使用 bool 查询组合多个条件
*/
public List complexBoolQuery(String keyword, List categoryIds,
Double minPrice, Double maxPrice,
List excludeStatus) {
EsQueryBuilder queryBuilder = EsQueryBuilder.builder(INDEX_NAME)
// must: 必须匹配(影响评分)
.match("title", keyword)
// filter: 必须匹配(不影响评分,可缓存)
.filterTerms("categoryId", categoryIds)
.filterRange("minPrice", minPrice, maxPrice)
// must_not: 必须不匹配
.mustNotTerms("status", excludeStatus)
// should: 可选匹配(有则加分)
.should(QueryBuilders.matchQuery("actor", keyword))
.minimumShouldMatch(0);
return esService.search(queryBuilder, Program.class);
}
/**
* 5. 使用 constant_score 包装不需要评分的查询
*/
public List constantScoreQuery(Long categoryId) {
// 使用原生 QueryBuilder 实现 constant_score
QueryBuilder query = QueryBuilders.constantScoreQuery(
QueryBuilders.termQuery("categoryId", categoryId)
).boost(1.0f);
return esService.searchByQueryBuilder(INDEX_NAME, query, 100, Program.class);
}
/**
* 6. 批量查询优化:使用 multi_search
*/
public Map> multiSearch(List categoryIds) {
// 注意:需要扩展 ElasticsearchService 支持 multi_search
// 这里展示思路
Map> result = new HashMap<>();
for (Long categoryId : categoryIds) {
EsQueryBuilder queryBuilder = EsQueryBuilder.builder(INDEX_NAME)
.filterTerm("categoryId", categoryId)
.fromSize(0, 10);
List programs = esService.search(queryBuilder, Program.class);
result.put(String.valueOf(categoryId), programs);
}
return result;
}
/**
* 7. 高亮优化:限制高亮片段
*/
public List> optimizedHighlightQuery(String keyword) {
EsQueryBuilder queryBuilder = EsQueryBuilder.builder(INDEX_NAME)
.match("title", keyword)
.match("description", keyword);
// 自定义高亮配置
HighlightBuilder highlightBuilder = new HighlightBuilder()
.field(new HighlightBuilder.Field("title")
.fragmentSize(100) // 片段大小
.numOfFragments(1)) // 片段数量
.field(new HighlightBuilder.Field("description")
.fragmentSize(150)
.numOfFragments(2))
.preTags("")
.postTags("");
queryBuilder.setHighlightBuilder(highlightBuilder);
return esService.searchWithHighlight(queryBuilder, Program.class);
}
}
```
#### **写入优化**
```java
/**
* 写入优化最佳实践
*/
@Slf4j
@Service
public class WriteOptimizationService {
@Autowired
private ElasticsearchDocumentService documentService;
@Autowired
private ElasticsearchIndexService indexService;
private static final String INDEX_NAME = "program";
/**
* 1. 批量写入优化
* - 合理设置批次大小(通常 1000-5000)
* - 控制并发写入线程数
*/
public void optimizedBatchWrite(List programs) {
int batchSize = 1000;
int totalSize = programs.size();
log.info("开始批量写入,总数: {}", totalSize);
long startTime = System.currentTimeMillis();
// 分批处理
for (int i = 0; i < totalSize; i += batchSize) {
int endIndex = Math.min(i + batchSize, totalSize);
List batch = programs.subList(i, endIndex);
BulkResult result = documentService.batchAdd(INDEX_NAME, batch);
if (result.hasFailures()) {
log.warn("批次 {}-{} 部分失败,失败数: {}",
i, endIndex, result.getFailureCount());
}
log.info("写入进度: {}/{}", endIndex, totalSize);
}
long elapsed = System.currentTimeMillis() - startTime;
log.info("批量写入完成,耗时: {}ms,速率: {} docs/s",
elapsed, totalSize * 1000L / elapsed);
}
/**
* 2. 大批量数据导入优化
* - 临时关闭副本
* - 增大 refresh_interval
* - 导入完成后恢复
*/
public void bulkImportWithOptimization(List programs) {
try {
// 优化索引设置
optimizeIndexForBulkImport();
// 执行批量导入
optimizedBatchWrite(programs);
} finally {
// 恢复索引设置
restoreIndexSettings();
}
}
private void optimizeIndexForBulkImport() {
Map settings = new HashMap<>();
settings.put("index.refresh_interval", "-1"); // 关闭自动刷新
settings.put("index.number_of_replicas", 0); // 临时关闭副本
indexService.updateSettings(INDEX_NAME, settings);
log.info("索引设置已优化");
}
private void restoreIndexSettings() {
Map settings = new HashMap<>();
settings.put("index.refresh_interval", "1s"); // 恢复刷新间隔
settings.put("index.number_of_replicas", 1); // 恢复副本
indexService.updateSettings(INDEX_NAME, settings);
// 强制刷新
indexService.refreshIndex(INDEX_NAME);
// 强制合并段
indexService.forceMerge(INDEX_NAME, 1);
log.info("索引设置已恢复");
}
/**
* 3. 使用异步写入
*/
@Async
public CompletableFuture asyncBatchWrite(List programs) {
BulkResult result = documentService.batchAdd(INDEX_NAME, programs);
return CompletableFuture.completedFuture(result);
}
/**
* 4. 并行批量写入
*/
public void parallelBatchWrite(List programs, int threadCount) {
int batchSize = programs.size() / threadCount;
ExecutorService executor = Executors.newFixedThreadPool(threadCount);
List> futures = new ArrayList<>();
try {
for (int i = 0; i < threadCount; i++) {
int start = i * batchSize;
int end = (i == threadCount - 1) ? programs.size() : start + batchSize;
List batch = programs.subList(start, end);
Future future = executor.submit(() ->
documentService.batchAdd(INDEX_NAME, batch));
futures.add(future);
}
// 等待所有任务完成
int totalSuccess = 0;
int totalFailure = 0;
for (Future future : futures) {
BulkResult result = future.get();
totalSuccess += result.getSuccessCount();
totalFailure += result.getFailureCount();
}
log.info("并行写入完成,成功: {}, 失败: {}", totalSuccess, totalFailure);
} catch (Exception e) {
log.error("并行写入失败", e);
throw new RuntimeException(e);
} finally {
executor.shutdown();
}
}
}
```
#### **缓存策略**
```java
/**
* ES 查询缓存策略
*/
@Slf4j
@Service
public class EsCacheService {
@Autowired
private ElasticsearchService esService;
@Autowired
private RedisTemplate redisTemplate;
private static final String CACHE_PREFIX = "es:program:";
private static final Duration CACHE_TTL = Duration.ofMinutes(5);
/**
* 1. 使用 Redis 缓存查询结果
*/
@SuppressWarnings("unchecked")
public List searchWithCache(String keyword) {
String cacheKey = CACHE_PREFIX + "search:" + DigestUtils.md5Hex(keyword);
// 先查缓存
Object cached = redisTemplate.opsForValue().get(cacheKey);
if (cached != null) {
log.debug("命中缓存: {}", cacheKey);
return (List) cached;
}
// 查询 ES
EsQueryBuilder queryBuilder = EsQueryBuilder.builder("program")
.match("title", keyword);
List result = esService.search(queryBuilder, Program.class);
// 写入缓存
if (!result.isEmpty()) {
redisTemplate.opsForValue().set(cacheKey, result, CACHE_TTL);
}
return result;
}
/**
* 2. 热门搜索词预热缓存
*/
@Scheduled(fixedRate = 300000) // 每5分钟
public void warmUpCache() {
List hotKeywords = getHotKeywords();
for (String keyword : hotKeywords) {
try {
searchWithCache(keyword);
} catch (Exception e) {
log.warn("预热缓存失败: {}", keyword, e);
}
}
log.info("缓存预热完成,预热 {} 个关键词", hotKeywords.size());
}
/**
* 3. 数据更新时清除相关缓存
*/
public void invalidateCache(Long programId) {
// 清除该节目相关的所有缓存
Set keys = redisTemplate.keys(CACHE_PREFIX + "*");
if (keys != null && !keys.isEmpty()) {
redisTemplate.delete(keys);
log.info("清除 {} 个缓存", keys.size());
}
}
/**
* 4. 分类统计结果缓存(变化不频繁的数据)
*/
@Cacheable(value = "es:stats", key = "'category'", unless = "#result == null")
public Map getCategoryStats() {
EsQueryBuilder queryBuilder = EsQueryBuilder.builder("program")
.termsAggregation("category_count", "categoryId", 100);
Map aggResult =
esService.searchAggregation(queryBuilder);
// 解析结果
Map stats = new HashMap<>();
AggregationResult categoryAgg = aggResult.get("category_count");
if (categoryAgg != null && categoryAgg.getBuckets() != null) {
for (AggregationResult.BucketData bucket : categoryAgg.getBuckets()) {
stats.put(bucket.getKey(), bucket.getDocCount());
}
}
return stats;
}
private List getHotKeywords() {
// 从统计服务获取热门搜索词
return Arrays.asList("演唱会", "话剧", "音乐剧", "相声", "脱口秀");
}
}
```
#### **监控与告警**
```java
/**
* ES 监控服务
*/
@Slf4j
@Service
public class EsMonitorService {
@Autowired
private RestHighLevelClient client;
@Autowired
private MeterRegistry meterRegistry;
/**
* 1. 记录查询耗时指标
*/
public List searchWithMetrics(EsQueryBuilder queryBuilder, Class clazz) {
Timer.Sample sample = Timer.start(meterRegistry);
try {
// 执行查询
List result = esService.search(queryBuilder, clazz);
// 记录成功指标
sample.stop(Timer.builder("es.query.duration")
.tag("index", queryBuilder.getIndexName())
.tag("status", "success")
.register(meterRegistry));
// 记录结果数量
meterRegistry.gauge("es.query.result.size",
Tags.of("index", queryBuilder.getIndexName()),
result.size());
return result;
} catch (Exception e) {
// 记录失败指标
sample.stop(Timer.builder("es.query.duration")
.tag("index", queryBuilder.getIndexName())
.tag("status", "error")
.register(meterRegistry));
meterRegistry.counter("es.query.error",
Tags.of("index", queryBuilder.getIndexName())).increment();
throw e;
}
}
/**
* 2. 集群健康检查
*/
@Scheduled(fixedRate = 60000)
public void checkClusterHealth() {
try {
ClusterHealthRequest request = new ClusterHealthRequest();
ClusterHealthResponse response = client.cluster()
.health(request, RequestOptions.DEFAULT);
String status = response.getStatus().name();
int numberOfNodes = response.getNumberOfNodes();
int activeShards = response.getActiveShards();
// 记录指标
meterRegistry.gauge("es.cluster.nodes", numberOfNodes);
meterRegistry.gauge("es.cluster.shards.active", activeShards);
// 状态告警
if ("RED".equals(status)) {
sendAlert("ES 集群状态异常: RED",
String.format("节点数: %d, 活跃分片: %d",
numberOfNodes, activeShards));
} else if ("YELLOW".equals(status)) {
log.warn("ES 集群状态: YELLOW,可能有副本未分配");
}
log.info("ES 集群健康检查: status={}, nodes={}, activeShards={}",
status, numberOfNodes, activeShards);
} catch (Exception e) {
log.error("ES 集群健康检查失败", e);
sendAlert("ES 集群健康检查失败", e.getMessage());
}
}
/**
* 3. 索引状态监控
*/
@Scheduled(fixedRate = 300000)
public void checkIndexStats() {
try {
IndicesStatsRequest request = new IndicesStatsRequest();
IndicesStatsResponse response = client.indices()
.stats(request, RequestOptions.DEFAULT);
response.getIndices().forEach((indexName, stats) -> {
long docCount = stats.getPrimaries().getDocs().getCount();
long storeSizeBytes = stats.getPrimaries().getStore().getSizeInBytes();
meterRegistry.gauge("es.index.docs.count",
Tags.of("index", indexName), docCount);
meterRegistry.gauge("es.index.store.size.bytes",
Tags.of("index", indexName), storeSizeBytes);
log.debug("索引 [{}] 统计: 文档数={}, 存储大小={}MB",
indexName, docCount, storeSizeBytes / 1024 / 1024);
});
} catch (Exception e) {
log.error("索引状态监控失败", e);
}
}
/**
* 4. 慢查询日志收集
*/
public void logSlowQuery(String indexName, String dsl, long elapsed) {
if (elapsed > 3000) { // 超过3秒
log.warn("[慢查询] 索引: {}, 耗时: {}ms, DSL: {}",
indexName, elapsed, dsl);
// 发送告警
if (elapsed > 10000) { // 超过10秒
sendAlert("ES 慢查询告警",
String.format("索引: %s, 耗时: %dms", indexName, elapsed));
}
}
}
private void sendAlert(String title, String content) {
// 发送告警通知(钉钉、企业微信、邮件等)
log.error("【告警】{}: {}", title, content);
}
}
```
## **十一、总结**
![[2-Learning/05-项目/08-企业级项目深读/02-damai_pro/05-业务深读/04-节目服务深度解析/assets/Elasticsearch_知识总结-3b295165.jpg]]
---
**企业级项目导航**:⬅️ [[10-单体项目加密脱敏最佳实践|10-单体项目加密脱敏最佳实践]] | 01-Elasticsearch 学习笔记 | ➡️ [[02-主页数据显示讲解|02-主页数据显示讲解]]